Part 2 · 2 chapters · ~20 min
Consensus
Raft step by step: terms, randomised timeouts, RequestVote and the majority that makes one leader per term, stale leaders stepping down; then log replication, commit on a majority, why committed entries survive, repairing followers, the cost, and what consensus buys.
5
Leader election
one leader per term, by majority vote
- Steady state: heartbeats reset randomised election timers.
- The leader fails, heartbeats stop, and the shortest timer fires first.
- The candidate increments the term, votes for itself and asks the others. Voters refuse a candidate whose log is behind theirs.
- A majority wins. Majorities overlap, so there is one leader per term. Split votes retry.
- A higher term always wins. Stale leaders step down.
- The returning old leader becomes a follower. It never had a majority, so it committed nothing.
| cluster size | majority | failures tolerated | note |
|---|---|---|---|
| 3 | 2 | 1 | the minimum for fault tolerance; one per zone |
| 4 | 3 | 1 | no better than 3, and more to coordinate |
| 5 | 3 | 2 | survives a zone plus one more server; the usual production choice |
| 7 | 4 | 3 | rarely worth the extra write latency |
where you already use it
Kubernetes stores all its state in etcd, which is Raft. CockroachDB, TiKV and YugabyteDB run a Raft group per range of data. Consul, Kafka's KRaft controller and many managed databases' failover coordinators are Raft or Paxos underneath. When the etcd cluster loses its majority, Kubernetes cannot schedule anything.
RAFT: LEADER ELECTION
five servers, randomised timeouts, terms, and the majority vote that keeps two leaders from ever coexisting in one term
swipe the figure sideways, or tap expand for full screen
1/6
steady state
Steady state, term 3: S1 is leader and sends heartbeats (empty AppendEntries) every 50 ms; followers S2 to S5 reset their election timers on every heartbeat. Each follower's timer is random in a range (150 to 300 ms) so they do not all time out together.
6
Log replication, commit, and what consensus buys
a write is safe once a majority has it
- The client writes to the leader, which appends the entry uncommitted.
- Replicate with AppendEntries, which includes a consistency check on the previous entry.
- Commit once a majority has stored the entry: apply it, answer the client, propagate the commit index.
- Why a majority: any two majorities overlap, so committed entries survive every change of leader.
- Repair a lagging follower by walking back to the last entry it shares with the leader.
- The cost: a majority round trip per write, a leader bottleneck, and no writes in a minority partition.
| consensus buys | example |
|---|---|
| a single agreed leader | exactly one primary per database shard; no split brain |
| an agreed order of operations | a replicated log every replica applies identically (state machine replication) |
| linearizable reads and writes | a configuration value or a lock that every client sees the same way |
| atomic membership changes | adding or removing servers without two overlapping configurations |
| fencing tokens | a monotonically increasing term or lease id that storage can check to reject a stale leader's writes |
Paxos
Lamport's Paxos (1989, published 1998) solved the same problem first and underlies Chubby and Spanner. Raft is equivalent in what it guarantees and much easier to implement correctly, which is why most systems built since 2014 chose it. Part 9 covers the papers.
RAFT: LOG REPLICATION AND COMMIT
how a write becomes durable on a majority, what "committed" means, and how a lagging follower is repaired
swipe the figure sideways, or tap expand for full screen
1/6
client write
The client writes: "set balance:acct_7 = 1500" arrives at the leader S3 (term 4). S3 appends it to its log at index 7, term 4, uncommitted. Followers that receive writes redirect clients to the leader.