Raft is usually taught as a list of rules: vote once per term, refuse candidates with stale logs, only commit entries from the current term. Each rule looks arbitrary until you see the failure it prevents. Then the algorithm becomes small enough to rebuild from memory, which is what you need when reading an incident timeline for etcd, Consul or your own replicated log.
This article derives Raft instead of listing it. It starts with the simplest replication scheme that could work, breaks it, and adds one mechanism per failure. The rules as a production node implements them, including reads, PreVote and membership changes, are covered in Raft consensus, in depth; this page is the intuition that makes those rules obvious.
The problem Raft solves
You want several servers to behave like one reliable machine. The standard way is a replicated state machine: every server runs the same deterministic program over the same sequence of commands, so every server ends in the same state. The hard part is agreeing on the sequence when servers crash, restart and lose messages.
The guarantee is simple to state: once a client is told a command is committed, that command stays at the same log position on every server forever. Raft assumes servers fail by stopping, not by lying, and that messages can be delayed, dropped or reordered. It never relies on timing for correctness, only for progress.
Attempt 1: a primary and a backup
Start with the obvious design. One server is the primary. Clients send commands to it, it appends them to its log, copies them to a backup and replies. If the primary dies, the backup takes over.
It breaks the first time the network, rather than a server, fails. Partition primary from backup: the backup stops hearing heartbeats and promotes itself while the primary keeps serving clients on its side. Two primaries now accept conflicting writes: split brain. Waiting longer does not help, because a dead primary and an unreachable one look identical. Any design in which one server can decide alone that it is in charge will eventually have two in charge.
Fix 1: decisions need a majority
The way out is to make every decision need a majority of the configured cluster: 3 of 5, 2 of 3. The key property is overlap: any two majorities share at least one server. Two leaders cannot win the same election, because the shared server votes once, and any future majority includes a server that has an acknowledged write.
The majority is of the configured membership, not of reachable servers; counting only reachable servers turns every partition back into split brain. The cost is availability: five servers tolerate two failures, three tolerate one, and a fourth server adds nothing, which is why clusters have odd sizes. Paxos reaches the same overlap argument from a different direction.
Fix 2: terms make old leaders harmless
Majorities stop two leaders winning one election, but not an old leader lingering. Leader A is partitioned away, the rest elect B, the partition heals, and A still believes it leads.
The fix is the term, an integer that increases with every election. A candidate increments its term and asks for votes in it, and every message carries the sender's term. A server that sees a higher term adopts it and becomes a follower; a message with a lower term is rejected. A's stale appends are refused, the replies carry B's term, and A steps down. A term is an epoch, the same idea as fencing tokens.
This is also why each server persists its term and vote before replying. A server that votes, crashes, forgets and votes again in the same term breaks the overlap argument.
Fix 3: a candidate must have every committed entry
Now there is at most one leader per term. The next failure is subtler. A write reaches S1, S2 and S3 of five and is committed; then S1, the leader, crashes. If S4, which never got the write, wins votes from S4, S5 and one more server, a committed command disappears.
Overlap saves us again with one more rule. Any electing majority includes a server holding the committed entry, so let each voter refuse a candidate whose log is less up to date than its own. Raft compares the last entry's term first and log length second. A server holding the committed entry refuses S4, which cannot assemble a majority. This election restriction means committed entries only flow from leaders to followers.
def log_ok(my_last_term, my_last_index, cand_last_term, cand_last_index):
"""The election restriction: grant a vote only to a log at least as up to date as mine."""
if cand_last_term != my_last_term:
return cand_last_term > my_last_term # later last term wins outright
return cand_last_index >= my_last_index # same last term: longer or equal log wins
def on_request_vote(me, req):
if req.term > me.current_term: # any newer term: step down, forget old vote
me.current_term, me.voted_for, me.role = req.term, None, "follower"
if req.term < me.current_term:
return False # stale candidate
if me.voted_for not in (None, req.candidate_id):
return False # one vote per term
if not log_ok(me.last_log_term(), me.last_log_index(),
req.last_log_term, req.last_log_index):
return False
me.voted_for = req.candidate_id
me.persist() # term and vote reach disk BEFORE replying
return TrueComparing last terms before lengths matters. A long log full of entries from a deposed leader's old term is less trustworthy than a shorter log that ends with an entry from a newer term, because the newer leader was elected with the election restriction already in force.
Fix 4: log matching keeps the logs honest
Replication needs one invariant. With each batch of entries the leader sends the index and term of the entry just before them. A follower accepts only if its log has an entry at that index with that term; otherwise it rejects and the leader steps back and retries. When the check passes, the follower drops conflicting entries after the match point and appends the leader's.
The effect is an induction argument: if two logs share an entry with the same index and term, they are identical up to that index, because every append checked its predecessor. So committing one index commits every earlier one, and the leader repairs a follower by finding the last point of agreement and overwriting from there. Followers can lose uncommitted entries this way, which is correct because no client was told they were committed.
Fix 5: the commit rule, and why most is not enough
This is the rule that surprises everyone. It seems natural that an entry is committed as soon as a majority of servers store it. The diagram shows why that is wrong, and it is the scenario from Figure 8 of the Raft paper by Diego Ongaro and John Ousterhout.
Here is how the cluster got there. In term 2, S1 led and copied index 2 to S2 only. S1 crashed; S5 won term 3 with votes from S3, S4 and itself, and wrote its own index-2 entry, but crashed before sending it anywhere. S1 restarted and won term 4, then copied its old term-2 entry to S3. Index 2 is now on three of five servers.
Suppose S1 treats that as committed and then crashes. S5 restarts and asks for votes in term 5. Its last entry has term 3, while S2, S3 and S4 end in term 2 or 1, so under the election restriction all three find S5 at least as up to date and grant their votes. S5 wins and, through log matching, overwrites index 2 on everyone with its term-3 entry. A committed command has vanished.
Raft's fix: count replicas only for entries of the leader's current term. S1 appends a term-4 entry at index 3 and replicates that. Once it is on a majority, S5 cannot win, because a majority's logs now end in term 4, and committing index 3 commits index 2 with it. Running the vote check confirms it: before the term-4 entry spreads, S2, S3 and S4 all accept S5; after it reaches S2 and S3, only S4 does. That is why a new leader immediately appends a no-op in its own term.
def advance_commit(log, match_index, commit_index, current_term, cluster_size):
"""log: list of entry terms (index 1 is log[0]). match_index: follower -> highest replicated index."""
for n in range(len(log), commit_index, -1):
if log[n - 1] != current_term:
continue # never count replicas for an older term's entry
replicas = 1 + sum(1 for m in match_index.values() if m >= n) # 1 = the leader itself
if replicas > cluster_size // 2:
return n # commits n and, by log matching, everything before it
return commit_indexIn the scenario, advance_commit([1, 2, 4], {'S2': 2, 'S3': 2, 'S4': 1, 'S5': 1}, 1, 4, 5) returns 1: nothing new is committed even though index 2 has three copies. Once S2 and S3 match index 3, the same call returns 3.
Timeouts: liveness, not safety
Nothing so far depends on clocks. Clocks only decide when a follower gives up on a silent leader. Raft randomises each server's election timeout (the paper's example range is 150 to 300 milliseconds) so that usually one server times out first and wins before others wake up. A split vote simply times out again and retries in a new term.
The paper's rule of thumb is an ordering: broadcast time much less than election timeout, much less than mean time between failures. A timeout shorter than a round trip plus a disk sync makes followers depose a healthy leader; a very long one means seconds of unavailability after a crash. Production systems add PreVote so a rejoining server with an inflated term cannot depose a healthy leader.
How the rules fit together
| Failure | Rule that prevents it | What it costs |
|---|---|---|
| Two primaries after a partition | Every decision needs a majority of the configured cluster | Minority side cannot serve writes |
| An old leader keeps writing | Terms: higher term wins, lower term rejected | Term and vote must be persisted per election |
| New leader missing a committed write | Election restriction: vote only for logs at least as up to date | Some servers can never win until they catch up |
| Divergent follower logs | Log matching check on every append | Repair walks back to the last agreement point |
| Majority-stored old entry overwritten | Commit only by counting current-term entries | A new leader must replicate a no-op before committing old entries |
| Endless split votes | Randomised election timeouts | Seconds of unavailability after a crash, tuned by timeout |
Failure modes you will meet in practice
- Skipping the fsync. Acknowledging a vote or append before it is durable lets a power loss erase it, silently voiding the majority argument.
- Unchecked leader reads. A deposed leader in a minority partition serves stale reads until it learns a higher term. Linearisable reads need a quorum round or a lease.
- Even-sized or stretched clusters. Four servers across two zones lose quorum when either zone fails; three zones survive a zone loss.
- LAN timeouts on a WAN. Cross-region round trips plus disk latency can exceed a default timeout and cause constant elections.
- Assuming committed means applied. Servers apply committed entries asynchronously, so a follower read right after a write may miss it.
A worked example to try
Take the five logs in the diagram. Who can win if S1 crashes now? S5, with S2, S3 and S4, because its last term 3 beats their 2 or 1. What must S1 do before index 2 is safe? Get a term-4 entry onto at least two followers. After index 3 reaches S2 and S3, S5 collects only its own vote and S4's, two of five, and cannot win.
Then run it: model each server as a list of terms, use the two functions above and replay the scenario with different crash points. That fifty-line simulator is how the numbers here were checked.
What to do next
- Rebuild the table above from memory: for each rule, name the failure it prevents.
- Implement log_ok and advance_commit, then replay the Figure 8 scenario and confirm that counting old-term entries loses a committed write.
- Read the handlers and read-path rules in the production-oriented Raft article and map each line to one of the fixes here.
- In your own cluster, check that the member count is odd, that members span failure domains, and that the storage layer really fsyncs before acknowledging.
- Measure heartbeat round-trip plus disk sync latency at the 99th percentile and confirm the election timeout sits well above it.
- Check whether your system uses PreVote and how it serves reads, since those decide behaviour during partitions.