computer science//distributed systems//distributed consensus//Raft

Raft is a distributed consensus algorithm that keeps a replicated log identical across a cluster of servers by electing one leader, which orders every write and commits it once a majority of servers has stored it. Diego Ongaro and John Ousterhout designed it in 2014 to be easier to understand and implement than Paxos, and it now runs under much of everyday infrastructure: etcd, the store that holds the state of every Kubernetes cluster, uses Raft, so fleet software deployed on Kubernetes already depends on it.


Raft is a distributed consensus algorithm that keeps a replicated log identical across a cluster of servers by electing one leader, which orders every write and commits it once a majority of servers has stored it. Diego Ongaro and John Ousterhout designed it in 2014 to be easier to understand and implement than Paxos, and it now runs under much of everyday infrastructure: etcd, the store that holds the state of every Kubernetes cluster, uses Raft, so fleet software deployed on Kubernetes already depends on it.

The algorithm has three parts. In leader election, every follower runs a timer with a randomized timeout; if it hears nothing from a leader before the timer expires, it becomes a candidate, increments the term number and asks for votes, and a candidate that gets a majority leads for that term (the randomness makes it unlikely that two candidates split the vote twice in a row). In log replication, clients send writes to the leader, which appends them to its log, forwards them to the followers and marks an entry committed once a majority holds it; only committed entries are applied to the replicated state (state machine replication). The safety rule closes the loop: a candidate whose log is missing a committed entry cannot win an election, so a new leader always holds everything committed before it.

With 2f+12f+12f+1 servers Raft survives fff crashes, and during a network partition only the side holding the majority keeps writing. Three servers tolerate one failure and five tolerate two. The minority side refuses writes instead of accepting ones that might conflict, which makes Raft the consistent choice under the CAP theorem.

Its timeouts are how it lives with FLP impossibility. It never commits two values whatever the timing, and it elects a leader and progresses only while messages arrive within the election timeout often enough. The timeout is set well above the network round trip (hundreds of milliseconds inside a datacentre) so that a slow link does not trigger elections constantly.

Every write costs a round trip from the leader to a majority and a durable write to disk on each, so a Raft cluster holds small, critical state (who owns which task, configuration, locks), never streaming telemetry.

Raft tolerates crashes only. A server that lies can corrupt it; untrusted participants need Byzantine fault tolerance.

Its leader orders writes, and nothing is averaged. Averaging drone positions is the job of the consensus protocol, which shares the word and none of the machinery, and the leader of leader-follower control is a different animal again.