HotShard
Open the simulatorSimulator β†’
Distributed consensus

Raft Consensus

How a cluster of servers agrees on one ordered log of writes, even when some of them crash.

Any system that copies its data onto several machines hits the same question: when the copies disagree, whose version wins? Raft answers it by electing one leader at a time and sending every change through that leader, so there is always a single agreed order. This page covers the two halves that make it work, electing the leader and replicating the log, and the two rules that keep committed data safe.

Open the simulator β†’~8 min read

The problem: when replicas disagree, who wins?#

TL;DRthe 30-second version
  • Consensus means getting several machines to agree on one ordered list of writes (the log), even when some crash or messages are lost.
  • Raft elects one leader per term. Every write goes through the leader, so there is only ever one order.
  • A write is committed once a majority of servers store it. A cluster of 2f+1 servers survives f crashes, because any two majorities overlap.
  • Two rules keep committed data safe: a server that is behind on its log can't win an election, and a new leader only counts its own entries toward commit.

A single database server is a single point of failure. If it dies, every request fails until it's replaced. So you copy the data onto several machines. But now if two clients write to two different copies at the same time, which write wins, and in what order does every copy apply them?

Picture etcd, the key-value store that holds all of Kubernetes's cluster state. Its data is just a map: set x = 1, set y = 2, read x. Instead of copying the map around, each server stores the ordered list of writes. That list is the log. Each server rebuilds its own map by replaying the log in order. Same log, replayed the same way, gives every server the identical map.

So the whole job is one thing: get every server to agree on a single ordered log, despite crashes and a flaky network. That agreement is what consensus means.

Raft's answer is to elect exactly one leader at a time. Every write goes through the leader. The leader appends it to its own log and copies the log, in order, to the other servers (the followers). One leader means one order. The hard question of "what happened in what order" becomes "whatever order the leader wrote things down in."

Leader election: choosing the one leader#

The leader can crash, so the other servers need a way to notice and pick a new one with no human involved.

Every follower keeps a countdown timer. Each time it hears from the leader (a short "I'm still here" message called a heartbeat) it resets the timer. If the timer runs out, the follower assumes the leader is gone and starts an election.

Each election happens in a numbered round called a term. Term numbers only ever go up, so the term tells every server which election is the most recent.

  1. A follower's timer runs out. It moves to the next term, becomes a candidate, and votes for itself.
  2. It asks every other server for a vote with a message called RequestVote, tagged with its new term.
  3. A server grants the vote if the candidate's term is at least as new as its own and it hasn't already voted in this term. Each server gets one vote per term.
  4. If the candidate collects votes from a majority (more than half the servers, counting its own) it becomes leader and sends a heartbeat at once, so no other timer runs out.
  5. If two followers time out together they can split the vote, and no one reaches a majority. Each waits a fresh random time (say 150 to 300 ms) and tries again, so the tie breaks quickly.
followerthe default state
timer runs out β†’
candidatewhen its timer runs out
majority of votes β†’
leaderwon a majority of votes
The three roles, and how a server moves between them

That footer is how a stale leader fixes itself. Every message carries the sender's term. A leader that was cut off by a network problem keeps believing it's in charge. The moment it hears a message with a higher term, it knows a newer election has happened, and it steps down to follower on its own. A higher term always wins.

PredictFive servers. Two of them time out at the same moment and both become candidates in term 4. Can both become leader?

Hint: How many votes does each need to win, and how many votes exist in total?

No. Each needs a majority, 3 of 5, and there are only 5 votes, each server casting at most one. Two candidates can't both collect 3 votes from the same 5 servers, so at most one wins term 4. If neither reaches 3 (a split vote), they wait a random time and retry in term 5. Majority plus one vote per term is exactly what makes at most one leader per term certain.

Log replication and the commit rule#

Once elected, the leader accepts client commands. It appends each one to its own log first, then sends it to the followers with a message called AppendEntries. This is the same message it uses as a heartbeat, just carrying entries instead of being empty. Each AppendEntries names the entry that should come just before the new one, so the follower can check its log lines up before accepting.

If the follower's log doesn't line up at that spot (it's missing entries, or has different ones from an old, abandoned leader), it rejects the message. The leader tries one position earlier next round, and keeps stepping back until it finds a point where the two logs agree. Then it overwrites everything after that point on the follower with its own entries. Logs are never merged. The leader's log is the truth, and followers converge to match it.

An entry is committed once the leader has copied it to a majority of the servers. From that moment it is safe forever, even if the leader crashes the next instant. Any future leader is also chosen by a majority, and two majorities of the same cluster always share at least one server, which still holds the entry. Only after commit does each server apply the entry to its map, and only then does the leader tell the client the write succeeded.

index β†’ entry, each tagged with the term it was created in Β· highlighted = committed

logx=11 Β· t1y=22 Β· t1x=33 Β· t2z=94 Β· t3y=75 Β· t3
The replicated log on the leader

The two rules that keep committed data safe#

Elections and replication are enough to keep the cluster running, but not quite enough to keep it correct. Two extra rules make sure a committed entry can never be lost, no matter how leaders come and go.

Rule 1: a server that is behind can't win an election. A server grants its vote only if the candidate's log is at least as up-to-date as its own (compared by the last entry's term, then its position). A committed entry is held by a majority, so any candidate missing it is turned down by that whole majority and can never reach the votes it needs. A new leader is always guaranteed to already hold every committed entry. The paper calls this the election restriction.

PredictFive servers; three of them hold committed entry #10. A candidate has the highest term in the cluster, but its own log ends at entry #7. Can it win?

Hint: A vote now needs more than a fresh term. What is the extra check?

No. The three servers holding entry #10 all turn this candidate down because its log is shorter, so it can never reach a majority of 3. At most it gets itself plus the one other server that is also behind, which is two votes. A high term lets you start an election; an up-to-date log decides who can win one.

Rule 2: a new leader only counts entries from its own term toward commit. Here's the subtlety interviewers love to probe. Suppose an old leader copied an entry to a majority but crashed before marking it committed, and the new leader inherited that entry alongside its own newer ones. The new leader may not call it committed just because a majority now holds it, because a future leader that never saw it could still overwrite it. Once one of the new leader's own entries commits, every earlier entry beneath it, including the inherited one, is committed too.

Put the rules together and a network split works itself out. If the leader is cut off from a majority, it keeps accepting writes into its own log, but it can never replicate them to a majority, so they never commit. The majority side stops hearing heartbeats, times out, and elects a new leader among themselves. When the split heals, the old leader sees the higher term and steps down. The consistency check then overwrites its stuck, uncommitted writes with the new leader's log.

Quorums, fault tolerance, and message cost#

Raft's safety rests on one counting fact: any two majorities of the same set overlap in at least one member. If an entry is on a majority, and a future leader was elected by a majority, those two majorities share a server. By the election restriction, that server forced the new leader to already hold the entry.

  • Fault tolerance: a cluster of N = 2f+1 servers tolerates f crashes and still has a majority of f+1. 3 servers tolerate 1 failure; 5 tolerate 2; 7 tolerate 3.
  • Why odd sizes: 4 servers need a majority of 3, which is the same fault tolerance as 3 servers but with one more server to wait on. Clusters are almost always 3, 5, or 7.
  • Commit latency: one write needs a round-trip from the leader to the fastest majority of followers, roughly one network round-trip (RTT) plus a disk fsync on each. A cluster spanning two coasts (about 70 ms RTT) can't commit faster than about 70 ms per write, no matter how fast the disks are. That is why voting members sit in nearby zones.
  • Message cost: each AppendEntries round is O(N) messages, and so is an election. That is why Raft groups stay small (3 to 7) and scale by sharding into many groups, not by adding voters.
Trade-offs: what consensus costs you

Raft buys you a single, durable, linearizable log (every read returns the latest committed value). You pay for it in latency, a throughput ceiling, and availability during a network split. The side of a split with no majority refuses writes to protect consistency. Choosing consistency over availability like this is what people mean when they call Raft a CP system, in the language of the CAP theorem.

  • Availability vs consistency: under a partition, only a majority side stays writable. A cluster with no majority side is fully unavailable for writes, by design.
  • Latency: every committed write pays a majority round-trip plus a per-server fsync. You cannot beat one RTT to your quorum.
  • Leader bottleneck: all writes funnel through one server. Its CPU, disk, and uplink cap write throughput, and adding voters doesn't raise it.
  • Small clusters only: past 3 to 7 servers, quorums get slower. Scaling data means sharding into many independent Raft groups, each with its own leader (the CockroachDB and TiKV model, often called multi-Raft).
When Raft is the wrong toolIf you don't need linearizability, say for a cache or analytics rollups, Raft's majority round-trip is pure overhead, and a leaderless, gossip-style design gives you higher availability and write throughput. Reach for Raft when you need a strongly consistent, ordered source of truth: metadata, configuration, locks, transaction status.
Raft vs Multi-Paxos vs ZAB
RaftMulti-PaxosZAB (ZooKeeper)
LeadershipStrong single leader; all entries flow leader to followerOptional distinguished proposer; leadership is a convention, not coreSingle leader (primary) elected per epoch
Log modelAppend-only log with no holes; followers mirror the leaderPer-slot agreement; logs can fill out of order with gapsBroadcasts entries in the primary's order
Primary goalUnderstandability and practical implementationMinimal, general consensus theoryTotal-order broadcast for a primary-backup store
Used byetcd, Consul, CockroachDB, TiKV, Kafka KRaftGoogle Chubby, Spanner, MegastoreApache ZooKeeper, and systems built on it (HBase, Kafka before KRaft)

Raft and Multi-Paxos are equally powerful. Raft is essentially Multi-Paxos with a strong-leader rule and a no-gaps log, added to make it teachable and buildable. ZAB is close to Raft in spirit; what Raft calls a term, ZAB calls an epoch.

Where Raft runs in the wild
  • etcd: the canonical Raft library and the key-value store behind Kubernetes. The entire cluster state of Kubernetes lives in a single etcd Raft group.
  • HashiCorp Consul and Nomad: service discovery and scheduling on the hashicorp/raft library.
  • CockroachDB: each shard (range) of the keyspace is its own Raft group, thousands per cluster.
  • TiKV / TiDB: the same per-region multi-Raft model. TiKV's Raft is a direct port of etcd's to Rust.
  • Kafka KRaft: the controller quorum runs a Raft-based metadata log, replacing the separate ZooKeeper ensemble.
Gotchas a real deployment has to handle
  • The log grows forever. Each server periodically saves a snapshot (a compact picture of its current state) and throws away the log entries it covers. A follower that has fallen too far behind is sent the whole snapshot instead of the missing entries.
  • Changing the set of servers is risky. If servers switch to the new set at different moments, the old set and the new set could each elect a leader for an instant. Change membership one server at a time: then the old majority and the new majority always share a server, so they can't split into two leaders.
  • A server cut off from the cluster keeps timing out and raising its term. When it reconnects, its inflated term makes the healthy leader step down for no good reason. Pre-vote fixes this: before raising its term, a candidate first asks the others whether they would even vote for it.
  • Raft does not tolerate Byzantine (lying or buggy) servers. It assumes crash-stop failures and a network that may lose, delay, or reorder messages but not forge them. That needs a different class of protocol (PBFT, Tendermint).
In an interview

Lead with the shape of the problem: replicate a log by funnelling all writes through a single elected leader per term. That turns "how do N machines agree" into "how do we safely elect and replace a leader", a much smaller problem. Then know the two messages by name (RequestVote, AppendEntries) and the two safety rules.

Can two servers ever both be leader and commit conflicting entries?

Two servers can briefly both believe they are leader, for example a partitioned old leader that hasn't heard about the new term. But they cannot both commit. Commit needs a majority, terms only increase, and each server votes once per term, so two leaders can't exist in the same term, and the stale leader can never reach its own majority. Its writes stay uncommitted and are later overwritten.

A 5-server cluster splits 3 to 2, with the old leader on the 2-server side. What happens?

The 3-server side has a majority, so it elects a new leader in a higher term and keeps committing. The 2-server side can never reach 3, so any writes the old leader accepts stay uncommitted. On heal, the old leader sees the higher term, steps down, and its uncommitted suffix is overwritten by the new leader's log. Logs are never merged.

Are reads automatically linearizable because writes go through Raft?

No, and this is the most common Raft mistake. A server that thinks it's leader might be a deposed stale leader serving old data. For a linearizable read, the leader first confirms it still leads with one heartbeat round to a majority, or holds a short time-limited lease that lets it skip the check. Reading from a follower, or from an unconfirmed leader, can return stale results.

Why are clusters 3, 5, or 7 servers?

A cluster of 2f+1 tolerates f crashes. Adding one more server to make it even raises the majority size without raising the failures it survives: 4 servers still tolerate only 1, the same as 3, with one more server to wait on.

References & further reading
References

Feedback on this topic β†’