Consensus & Raft
Five machines keep one history of a shop's stock: elections and terms, commit on a majority, what crashes and partitions do, and why clusters come in odd sizes.
An interactive System Design lesson: 19 steps, about 32 minutes, on a live simulation in your browser.
A shop keeps its stock counts on five machines, node-a to node-e, so that losing one or two loses nothing. There is one rule they must never break: all five apply the same changes in the same order. If two of them each decide on their own who bought the last lamp, it has been sold twice.
Getting a group to agree on one sequence of decisions while members fail is consensus. It is hard for two reasons. A node cannot tell a slow peer from a dead one: both are silent. And whatever it does about that silence, two nodes must never both decide.
What you will learn
Why agreement is hard
- Five machines, one history: Consensus is getting a group to apply the same decisions in the same order, even though no member can tell a slow peer from a dead one.
Electing a leader
- Who becomes the leader?: A node that hears no leader for its election timeout becomes a candidate; a candidate with votes from a majority becomes leader.
- Terms: a clock made of elections: A term is a logical clock: higher terms win, lower terms are ignored, and one vote per node per term means at most one leader per term.
- Why the timeouts are random: Random election timeouts make one node fire first. Equal timeouts would make every election a tie.
Replicating a write
- When is the write safe?: An entry is committed once a majority stores it. From then on every possible future leader has it, so the leader may answer the client.
- Committed, then applied: The log is the truth; the state machine is the log applied in order. A node applies an entry only once it knows the entry is committed.
Crashes
- Break it: a follower crashes: A Raft group keeps committing as long as a majority is up and connected. A crashed follower costs nothing until too many join it.
- The leader crashes too: When the leader goes silent, a follower's timer fires and it runs in a higher term. A split vote only costs another round.
- The old leader returns as a follower
Partitions
- The leader ends up in the minority: During a partition only the side with a majority can elect a leader and commit. A leader cut off in the minority keeps its title and can do nothing with it.
- Two leaders, one that counts: Two nodes may believe they lead, in different terms. Only the one that can reach a majority can commit, so there is still one history.
- Healing: the stale leader steps down: A stale leader steps down the moment it sees a higher term. Its uncommitted entries are overwritten; nothing that was committed is ever lost.
Rules, sizes and real systems
- A node with an old log runs for leader: The election restriction: nobody votes for a candidate whose log is behind theirs, so a leader always holds every committed entry.
- Break it: two down, then three: N nodes need floor(N/2) + 1 to agree and survive floor((N - 1)/2) failures. An even-sized cluster tolerates no more failures than the odd size below it.
- Drill: majority arithmetic
- Drill: failures a cluster survives
- Where Raft runs
Recap & playground
- Cheat sheet
- Playground