The problem: one log on five machines

A replicated state machine runs the same program on several servers and feeds every copy the same commands in the same order. If the program is deterministic, all copies end in the same state, and the service survives as long as enough copies are alive. The hard part is agreeing on the order while servers crash and messages are lost or delayed. Raft (Ongaro and Ousterhout, 2014) solves it with a strong leader: one server decides the order and the others copy its log. With 5 servers, any 3 form a majority, so the cluster keeps working with 2 of them down, and any two majorities share at least one server. That overlap is what every safety argument below rests on.

On the pageMeaning
circles S1 to S5servers on a ring; colour = role (grey follower, orange candidate, light orange pre-vote, green leader, dark grey down); the bar under a circle is its election timer
log tableone row per server: role, term, vote, commitIndex, lastApplied, the state machine (x, y) and the log cells (index on top, term inside, command below); grey text = not committed on that server
next / matchthe leader's nextIndex and matchIndex for each follower
appthe client; its n-th write is x=n for odd n and y=n for even n
panelsthe recent elections with every vote and the reason for each refusal, and the leader's commit rule for its last entries

Every message takes one tick of the clock. An action runs the clock until nothing is left to do; while the cluster is quiet, time stands still and heartbeats are not simulated (press Step Tick to see one round).

Terms: a logical clock

Time is divided into numbered terms. A term starts with an election and has at most one leader; a term can also end without a leader (a split vote). Every server stores its currentTerm and puts it in every message. A server that sees a higher term adopts it and becomes a follower at once, whatever it was; a message with a lower term is rejected. So a leader that was cut off and comes back learns from the first answer it gets that it has been replaced.

Leader election

A follower expects to hear from a leader regularly: the leader sends AppendEntries to everyone at least every heartbeat interval (2 ticks here), with no entries if there is nothing new. A follower that hears nothing for its election timeout becomes a candidate: it increments its term, votes for itself and sends RequestVote(term, lastLogIndex, lastLogTerm) to all others. Each server gives one vote per term, first come first served, and stores it on disk before answering (votedFor). A candidate with votes from a majority (3 of 5, itself included) is the leader and immediately sends AppendEntries so the others stop their own elections.

State diagram of a Raft server: a Follower becomes a Candidate on election timeout (term + 1, votes for itself); a Candidate becomes Leader with votes from a majority (3 of 5), starts again in the next term after a split vote, and returns to Follower when it sees a leader or a higher term; a Leader returns to Follower when it sees a higher term
A follower that stops hearing a leader runs for election; a majority of votes makes it leader, and any higher term sends it back to follower.

If two candidates start together, the votes can split: 2 and 2, nobody wins, and each candidate waits for a new timeout and tries again in the next term. Raft makes this unlikely by choosing every election timeout at random (the paper suggests 150 to 300 ms): usually one server times out well before the others and wins before anyone else starts. The page uses a fixed table of timeouts per server so every run is the same. Demo 4 sets two equal timeouts and gets a split vote; Demo 5 sets all timeouts equal: every election splits until the timeouts differ again.

Log replication

The leader turns every client command into a log entry with its index and the current term, appends it to its own log and sends it to the followers in AppendEntries(term, prevLogIndex, prevLogTerm, entries[], leaderCommit). A follower accepts only if its own log has an entry at prevLogIndex with term prevLogTerm: the consistency check. By induction this gives the log matching property: if two logs have an entry with the same index and term, they are identical up to that entry.

When the leader knows (from the followers' answers, kept in matchIndex) that an entry of its current term is stored on a majority, the entry is committed: the leader advances its commitIndex, applies the command to its state machine and only then answers the client. The next AppendEntries carries the new leaderCommit, and the followers apply the entry too. A committed entry is durable: it is on a majority, and (next section) every future leader will have it.

Logs of five servers, entries labelled with term and command. Index 1 to 3 (t1 x=1, t1 y=2, t2 x=3) are on S1, S2 and S3, so commitIndex is 3; index 4 (t2 y=4) is only on S1 and S2, so it is not committed; S4 has 2 entries and S5 is down
Entry 3 is on 3 of 5 servers, so the leader commits it and answers the client; entry 4 waits until a third server has it.

Repairing followers

A follower that was down or cut off, or one that holds entries of an old leader that never committed, fails the consistency check. The leader then lowers that follower's nextIndex by one and tries again, until it finds the last index where both logs agree. From there the follower deletes its conflicting entries and takes the leader's. Logs are never merged: the leader's log always wins, and the leader never deletes or overwrites anything in its own log. One step per round trip is the paper's basic rule; real implementations send a hint with the rejection (the conflicting term and its first index) so the leader can skip a whole term at once. Demo 6 shows three failed checks before S4 is repaired.

Safety: who may become leader, and what may be committed

Election restriction. A voter refuses a candidate whose log is less up to date than its own: the candidate's last entry must have a higher term, or the same term and at least the same index. Since a committed entry is on a majority and a winner needs a majority of votes, at least one voter has the entry and refuses anyone without it. So every leader has every committed entry (leader completeness).

Commit only entries of your own term. Counting copies is not enough for an entry of an older term. In the paper's Figure 8 (Demo 9), S1 leads term 2 and copies index 2 to S2 only; S5 wins term 3 without it and writes its own index 2; later S1 leads term 4 and copies its index 2 to a majority. If S1 now counted that as committed, S5 could still win an election (its last term 3 beats their 2) and overwrite index 2 everywhere: a committed entry would be lost. So a leader only counts replicas for entries of its current term; when one of those commits, everything before it commits too (Demo 10: once index 3 of term 4 is on a majority, S5 can no longer win). This is also why a new leader in real systems appends a no-op entry right after its election.

Leader crashes and clients

A client sends its command to the leader; a follower answers "not leader" and names the leader it knows. When the leader crashes, the followers notice only through missing heartbeats; after one election timeout a new leader is elected (Demo 3). A command that the old leader had appended but not committed may be lost or may be committed by the new leader: the client that timed out does not know which. If it simply sends the command again, it may be executed twice. Real Raft clients therefore give every command a client id and a sequence number, and the servers keep a client session table to drop duplicates. The page's client never resends a timed-out command, and shows at the end whether it was discarded or applied anyway.

Network partitions and reads

A leader cut off in a minority does not know it: nothing tells it, and in the basic algorithm it keeps its title. It still accepts commands, but it can never get them onto a majority, so they never commit and the clients time out; meanwhile the majority side elects a new leader in a higher term and carries on. When the partition heals, the old leader sees the higher term, steps down, and its uncommitted entries are deleted (Demo 8). Nothing that was acknowledged is lost, unlike asynchronous failover in Redis Sentinel or Redis Cluster.

Reads have the same trap: the old leader's state machine is stale, and answering from it would return data that the other side has already overwritten. With ReadIndex, the leader notes its commitIndex, sends one heartbeat round and answers only after a majority has replied, which proves nobody newer is leader; then it waits until it has applied up to that index. A new leader first needs a committed entry of its own term to know the latest commitIndex (hence the no-op). A lease read skips the round by trusting that no election can finish within the election timeout after the last majority heartbeat, which depends on bounded clock drift. CheckQuorum (etcd, the thesis) makes a leader step down by itself when it has not heard from a majority for an election timeout.

How this strong (linearizable) read compares with eventual, read-your-writes and causal reads: Consistency Patterns.

Pre-vote and disruptive servers

A server that is cut off alone keeps timing out and incrementing its term. When it comes back, its high term forces the working leader to step down, and the cluster holds a needless election that the returning server cannot even win, because its log is behind (Demo 6). With pre-vote (Ongaro's thesis, section 9.6) a server whose timer runs out first asks PreVote(term + 1) without changing its term; others agree only if they have not heard from a leader recently and its log is up to date. Only with a majority of pre-votes does it start the real election. The isolated server's term stays where it was, and after the heal it simply follows the leader (Demo 7).

What "persistent" means

currentTerm, votedFor and the log must be on stable storage (written and fsynced) before the server answers the message that changed them. Otherwise a server could crash, restart with an older term or no vote, and vote twice in one term; or acknowledge an entry it then forgets, letting the leader count a replica that does not exist. commitIndex, lastApplied and the state machine can be rebuilt: after a restart the page's servers start with commitIndex 0 and replay the log once a leader tells them what is committed.

Where Raft is used

etcd (and so every Kubernetes cluster's state), Consul, CockroachDB and TiKV (one Raft group per range of keys, thousands per node), and Kafka in KRaft mode, where the controllers keep the cluster metadata in a Raft-based log (Kafka page). The Redis Sentinel leader vote (Sentinel page) and the Redis Cluster replica election (Cluster page) use the same one-vote-per-epoch majority idea, without a replicated log.

Why a majority store refuses requests on the minority side, and what a leaderless Dynamo-style store does instead: The CAP Theorem: CP vs AP.

What the page leaves out

Snapshots and log compaction (InstallSnapshot), membership changes (joint consensus or one server at a time), batching and pipelining limits, learners (non-voting members), leadership transfer (TimeoutNow sent by the old leader; the page's Timeout Now button only fires one timer), disk and fsync latency, and flow control. Simplifications: the election timeouts come from a fixed table instead of fresh random numbers; messages that arrive in the same tick are handled nearest sender first; a server that cannot win stops after 2 elections per action; heartbeats are not simulated while the cluster is quiet; the leader's no-op is appended only when a read needs it; the client reaches every server even when the servers are partitioned; the log window shows the last 12 entries.