Consensus
Five servers have to agree on one list of writes while any of them can be cut off at any moment. Raft elects one of them to speak for a term and has the others copy it. Most of its rules read the way you would expect. One needed a figure of its own in the paper: a leader may count copies only of an entry from its own term. Leave that clause out and an entry that sat on a majority, and was acted on, can be deleted. This page lets you press every step of the paper's example of it, and on load tries every history three servers can reach in four terms, with one entry per log, to show that within that box the clause is the difference.
New to replicated logs? Start here
Several servers keep copies of one list of writes, so that losing a server loses nothing that was committed. The hard part is keeping the copies the same when any server can stop, or be cut off, between one message and the next.
Raft's answer is to elect one server to decide the order of writes for a while, called a term, and have the rest copy it. A write counts as committed once enough copies exist that no later leader can be elected without it. A majority is needed, and why is Majority's subject. This page is about why a majority of copies is not always enough.
A message that takes time and may not arrive
A message takes time to travel from one computer to another. Sending it does not establish that it arrived, and receiving one message does not establish that an earlier one arrived first.
Messages may be delayed, lost, duplicated or reordered. A protocol defines how the participants respond to those possibilities. Some machines here model only one of them; the page says which assumptions its result needs.
The machine for this idea on its own is Packet Switching, if you would rather press it than read about it.
Machines here that come first: Logical Clock.
Raft's election and replication, and the clause in its commit rule that names the term
1 Five servers or three, no leader, and the first timer to run out
Every server starts as a follower, waiting to hear from a leader. When its election timeout runs out and it has heard nothing, it adds one to its term, votes for itself and becomes a candidate. The paper draws each timeout at random, “chosen randomly from a fixed interval (e.g., 150–300ms)”, so that one server usually times out alone. The term is how servers tell old news from new; the paper says terms “act as a logical clock”, the kind Logical Clock builds. There is no clock on this page. You decide whose timer runs out, and when.
| server | role | term | voted for, this term | network |
|---|---|---|---|---|
| S1 | follower | 0 | nobody | all reachable |
| S2 | follower | 0 | nobody | all reachable |
| S3 | follower | 0 | nobody | all reachable |
| S4 | follower | 0 | nobody | all reachable |
| S5 | follower | 0 | nobody | all reachable |
Five followers in term 0, and no leader. Run a timer out.
Every step so far, oldest first:
- Nothing yet. Every server is a follower in term 0 with an empty log.
2 The candidate asks for votes, and what a voter refuses
A candidate asks every server it can reach for a vote. A server votes at most once in a term, for the first candidate to ask whose log is not behind its own. Refusing a candidate does not use up its vote. The paper: “the voter denies its vote if its own log is more up-to-date than that of the candidate”. Behind means the voter's last entry is from a later term, or both logs end in the same term and the voter's is longer. Votes from a majority of the whole cluster, counting its own, make the candidate the leader of its term. Why a second candidate cannot also collect a majority in that term is the subject of Majority, and it is not argued again here.
| server | its term | its last entry | voted for | its answer |
|---|---|---|---|---|
| Nobody is standing. Run a timer out in the first step. | ||||
- votes
- no candidate
Nobody is standing.
3 The leader copies its log out, and when an entry counts as committed
A leader sends each follower the entries it is missing, with the index and term of the entry just before them. The follower refuses unless it holds that entry, so the leader steps back one entry at a time until they agree. From there the follower deletes anything that conflicts and takes the leader's entries. Here that exchange is one step.
When does an entry count as committed, safe for every server to act on? Figure 2 of the paper: “If there exists an N such that N > commitIndex, a majority of matchIndex[i] ≥ N, and log[N].term == currentTerm: set commitIndex = N”. The last clause is what this page is about, and you can leave it out here.
Every server's log. Each box is one entry, holding the term it was written in; a heavy border marks a committed entry.
| server | 1 |
|---|---|
| S1, term 0, commit 0 | |
| S2, term 0, commit 0 | |
| S3, term 0, commit 0 | |
| S4, term 0, commit 0 | |
| S5, term 0, commit 0 |
The leader's count, entry by entry:
| index | term | held by | leader knows | counts | committed |
|---|---|---|---|---|---|
| Nobody leads, so nobody is counting copies. | |||||
Nobody leads, so no entry is being counted.
All five properties of the paper's Figure 3 hold in every state of this history, 0 steps.
4 Cut the network, heal it, and what the cut-off side deletes
Cut the network and each side carries on with the servers it can reach. The side with a majority can elect a leader and commit. A leader on the other side is not told it has been cut off. It keeps the title and keeps taking writes, and it can never commit them, because it can never hear from a majority. Heal the network and the first message across carries the newer term: the old leader steps down, and when the new leader sends it entries, it deletes every entry of its own that conflicts. Compare Two-Phase Commit, where a coordinator lost at the wrong moment leaves everyone waiting: here the side with a majority elects someone else and carries on.
On the far side of the cut:
- sides
- one network, all 5 servers reach each other
Every entry deleted so far:
| step | server | index | term | most servers that held it | had been committed |
|---|---|---|---|---|---|
| Nothing has been deleted. | |||||
Nothing has been deleted yet. Cut the network, let each side elect, write on both, and heal it.
5 An entry on a majority, deleted anyway: Figure 8, and every short history of three servers
The paper's Figure 8 is the last two steps together. An entry from term 2 reaches three of five servers and is deleted anyway, by a leader that never had it. Under Raft's rule nothing had counted it as committed, so nothing breaks. Count copies of any entry and the same history deletes an entry a server had already acted on. The paper says crashes; here a crashed server is one cut off from the rest, and it goes on believing it leads until a reply with a newer term reaches it.
Press a panel to replay the figure from an empty cluster, under the rule chosen above.
A figure is one history. To say the clause is what makes the difference, every history has to be tried. This page tries them when it loads, for three servers: every order of every timeout, vote, write and send, through four terms with one entry per log. In every state it reaches it checks the five properties the paper's Figure 3 says hold at all times: Election Safety, Leader Append-Only, Log Matching, Leader Completeness and State Machine Safety.
- with Raft's rule
- 15,675 states, nothing broken
- counting copies of any entry
- 14,424 states, broken after 13 steps
Both searches start from the same empty cluster and try every order of every timeout, vote, write and send. With Raft's rule nothing breaks. Without its last clause the search stops at its first broken state, 13 steps in, where S2 leads term 4 without the entry committed at index 1 in term 3.
All five properties of the paper's Figure 3 hold in every state of this history, 0 steps.
Your turn.
With Raft's rule on, get one entry onto a majority of the servers and then deleted, without breaking any of the five properties.
Not yet: nothing has been deleted.
These ran in this browser when the page loaded. Each claim, whether it held, and the number behind it.
| claim | held | measured |
|---|---|---|
| Election Safety holds in all 15,675 states three servers can reach, terms up to 4 and one entry per log | yes | 15,604 of those states had something for it to check |
| Leader Append-Only holds in all 15,675 states three servers can reach, terms up to 4 and one entry per log | yes | 26,691 appends checked |
| Log Matching holds in all 15,675 states three servers can reach, terms up to 4 and one entry per log | yes | 13,408 of those states had something for it to check |
| Leader Completeness holds in all 15,675 states three servers can reach, terms up to 4 and one entry per log | yes | 7,159 of those states had something for it to check |
| State Machine Safety holds in all 15,675 states three servers can reach, terms up to 4 and one entry per log | yes | 12,639 of those states had something for it to check |
| with the term check removed from the commit rule, the same search finds a history that breaks it, 13 steps long | yes | Leader Completeness fail: S2 leads term 4 without the entry committed at index 1 in term 3 |
| replayed with Raft's rule, the same history breaks nothing | yes | nothing was committed: the one entry on a majority came from an earlier term than its leader |
| with terms only up to 3, the other rule breaks nothing in 2,449 states, so its failure needs a fourth term | yes | the shortest failure ends in term 4 |
| the five panels of the paper's Figure 8 are reached exactly, by the same four actions the buttons press | yes | every log and every leader as printed, and nothing broken on the way |
| counting replicas of any entry, the same history breaks Leader Completeness and State Machine Safety on the way to (d) | yes | first at step 36 of 42: S5 leads term 5 without the entry committed at index 2 in term 4 |
| in (c) the term-2 entry is on 3 of 5 servers and the leader has not committed it; in (e) it is committed by the entry from term 4 after it | yes | S1's commit index is 1 in (c) and 3 in (e) |
| and from (e), S5 cannot win: it gets 2 votes of the 3 it needs | yes | only S4 and itself; S2 and S3 hold an entry from term 4 and S5's log ends in term 3 |
What is real here, and what is not
There is no clock, and you are the timer
Raft's elections run on timeouts, and the paper's randomised 150 to 300 milliseconds is why one server usually stands alone. This page measures no time at all. A timeout happens when you press it, in the order you press it. A split vote is something you can make on purpose, by running two timers out before either candidate asks anyone, and nothing here tells you how likely it would be on a real network.
A message arrives whole, with its reply, or not at all
Each vote request and each send is one step: delivered, answered and the answer read, before anything else happens. No message is delayed, duplicated or delivered out of order, and none is lost except by the cut you make. The one exception is deliberate: a vote request may reach a server after the candidate has already won, because a candidate asks everyone at once and the paper's Figure 8 counts a vote of that kind. The checker explores every order of these steps and every choice to never deliver one, which covers every cut. It does not explore a reply that arrives after its sender has moved to a newer term.
A crash is drawn as a cut
The paper's Figure 8 says servers crash and restart. A restarted server keeps its term, its vote and its log and loses the rest, including the fact that it was leader. Here a crashed server is cut off instead, and it keeps everything, so it goes on believing it leads. That is why the history for (c) and (d) has one more step than the caption: the old leader has to hear a newer term before it can stand again. Restarts are not in the checker either.
The leader's search for where two logs agree is one step
In the paper the leader keeps a nextIndex for each follower and steps it back one entry per refused message. Here the leader finds the last index where both logs hold an entry of the same term and sends everything after it, in one step. The end state is the same one the paper's loop reaches, and an independent implementation that does run the loop, message by message, is what the test suite compares this against.
Committing and applying are the same moment here
There is no state machine, no client and no reply to a client. An entry is taken to be applied the moment a server knows it is committed, which is the earliest the paper allows. State Machine Safety is checked on that: once a server has applied an entry at an index, no server may apply a different one there.
The checker's box is small, and the page says how small
Three servers, terms up to four, one entry per log. Within that it tries every history, and it counts two states as one when they differ only in which server is called which. Four terms is not a guess: with three, the rule without the clause breaks nothing either, and the page measures that on load. None of this is a proof for every cluster size. The paper reports a formal specification in TLA+, a machine-checked proof of one property and a written proof of State Machine Safety, and the specification's commit rule carries the same clause this page lets you remove.
Figure 8 is reached, and some of it is chosen
The figure does not say who led term 1 or who voted in (a), so the history here picks S2 for term 1 and S2 and S3 for (a). Everything the caption does say is held to: S5 wins term 3 with S3, S4 and itself, and in (d) with S2, S3 and S4. S5 also stands once in term 2 and loses, because S3 voted for S1 in that term; the figure does not show it, and nothing in the caption rules it out.
Not here: membership changes, snapshots, clients
The paper and its extended version cover changing the set of servers, compacting the log, and how clients find the leader. None of that is modelled. The page is leader election and log replication under a partition, and nothing else.
Why two majorities cannot disagree is on another page
Election Safety rests on two facts: a server votes once per term, and any two majorities of one cluster share a server. The second is Majority's whole subject, counted there for every pair, and this page uses it without proving it again.
Not the first consensus algorithm, and it does not claim to be
The paper sets itself beside Paxos and names Oki and Liskov's Viewstamped Replication as the algorithm it most resembles. Its claim is that Raft is easier to understand, and its evidence is a user study. That is why the page is dated to Raft and not to the problem.
Dates
The paper is in the proceedings of the 2014 USENIX Annual Technical Conference, June 19 to 20, 2014, in Philadelphia; the conference's own pages give June 17 to 20, which is the whole federated week around it. An extended version, with the sections on clients and log compaction that the conference version leaves out, is on the authors' site. Ongaro's dissertation, dated August 2014, carries the same figure as Figure 3.7 with its last two panels named (d1) and (d2).
Sound: no
Asked and answered. Nothing here has a duration: the page has no clock, and a sound for a vote or a deleted entry would decorate a step that the tables already print.
Sources
- D. Ongaro and J. Ousterhout, In Search of an Understandable Consensus Algorithm, 2014 USENIX Annual Technical Conference, June 19–20, 2014, Philadelphia. Figure 2's rules, Figure 3's five properties, Figure 8, and the 150 to 300 millisecond timeouts.
- The same paper, extended version, from the authors' site, with the same Figure 8 and caption and the sections on clients and log compaction.
- D. Ongaro, Consensus: Bridging Theory and Practice, PhD dissertation, Stanford University, August 2014. Figure 3.7 is the paper's Figure 8.
- D. Ongaro, the TLA+ specification of Raft. Its AdvanceCommitIndex carries the same clause: the entry at the new commit index must be from the leader's current term.
- ;login:, October 2014, the conference reports for USENIX ATC '14. The dates of the conference, from a second document.
- The Raft site, with RaftScope, the authors' own visualisation. Its commit rule keeps the clause, and nothing in it turns the clause off.
- Logical Art, the studio this belongs to.