How Raft consensus works
Raft explained for system design interviews: replicated state machines, leader election with terms and randomised timeouts, log replication and commit rules, safety, membership changes, snapshots, read consistency and where Raft is used in practice.
Reading is half of it. See this used in a real interview: walk through Design a Distributed Key-Value Store →
Whenever a design needs a group of machines to agree (who is the leader, what the configuration is, which lock holder is current), the answer underneath is a consensus algorithm, and today that usually means Raft. etcd, Consul, CockroachDB, TiKV and Kafka's KRaft controller all use it. You will rarely implement it in an interview, but explaining how it works shows you understand what "strongly consistent" really costs.
The goal: a replicated state machine
Every node keeps a log of commands and applies them in order to a state machine (for example a key-value map). If all nodes apply the same commands in the same order, they end up in the same state. Raft's job is to make every node agree on the log, despite crashes and network problems, as long as a majority of nodes is working. With 5 nodes, 2 can fail; with 3, 1 can.
Roles and terms
Each node is a follower, candidate or leader. Time is divided into numbered terms; each term has at most one leader. Terms act as a logical clock: messages from an older term are rejected, which fences out stale leaders. See distributed locks and leases.
Leader election
- Followers expect regular heartbeats from the leader.
- If a follower hears nothing for its election timeout (randomised, for example 150 to 300 ms), it becomes a candidate, increments the term, votes for itself and asks others for votes.
- Each node votes for at most one candidate per term, and only for a candidate whose log is at least as up to date as its own.
- A candidate with votes from a majority becomes leader and starts sending heartbeats.
Randomised timeouts make split votes rare: usually one node times out first and wins before others start.
Log replication
- Clients send commands to the leader.
- The leader appends the command to its log and sends it to followers (AppendEntries).
- Once a majority has stored the entry, the leader marks it committed, applies it, and responds to the client.
- Followers learn the commit index in later messages and apply committed entries.
Each AppendEntries message includes the index and term of the previous entry; a follower rejects it if its log does not match there, and the leader backs up until the logs agree, then overwrites the follower's divergent tail. This keeps all logs identical up to the commit point.
Safety
The election rule (vote only for candidates with up-to-date logs) guarantees that any new leader already has every committed entry, so committed entries are never lost or changed. A partitioned old leader cannot commit anything, because it cannot reach a majority; when it reconnects and sees a higher term, it steps down.
Costs
- Every write needs a majority round trip: latency is at least one network round trip to the nearest majority, plus disk syncs. Across regions, that is tens to hundreds of milliseconds.
- Throughput is bounded by the leader: all writes go through it. Systems scale by running many independent Raft groups, one per data range (multi-Raft, as in CockroachDB and TiKV). See sharding and partitioning.
- Majority required: lose a majority and the group stops accepting writes. That is the price of consistency. See CAP theorem.
Reads
Reading from the leader's state can still be stale if a newer leader exists that this node does not know about. Options:
- Run reads through the log (slow but simple).
- ReadIndex: the leader confirms it is still leader with a heartbeat round before serving the read.
- Leader leases: the leader serves reads locally while its lease is valid, relying on bounded clock drift.
- Follower reads for stale-tolerant queries.
Membership changes and snapshots
- Adding or removing nodes changes what "majority" means, so Raft changes membership carefully (one node at a time, or joint consensus) to avoid two majorities.
- Logs grow forever, so nodes periodically take a snapshot of the state machine and discard log entries before it. Lagging or new nodes receive the snapshot, then the remaining log.
Raft and Paxos
Paxos came first and underlies systems like Chubby and Spanner; Raft was designed to be easier to understand and implement, with a strong leader and a clear log structure. They provide the same guarantees. See consensus and coordination.
In the interview
You might say: "Metadata (partition assignments, leader of each shard) lives in a Raft-replicated store such as etcd, with three or five nodes. Writes commit on a majority, so we tolerate one or two failures; terms fence old leaders." For a strongly consistent key-value store, mention multi-Raft per range and the cross-region latency cost. For a job scheduler, leader election via the same store.
Checklist
- Replicated log applied in order to a state machine.
- Terms, randomised election timeouts, majority votes with log freshness.
- Commit on majority replication; logs repaired to match the leader.
- Majority round trip per write; many groups to scale.
- Linearizable reads via ReadIndex or leases; follower reads when staleness is fine.
- Careful membership changes; snapshots to truncate logs.