Consensus with Raft

How a cluster of unreliable machines agrees on one truth: Raft leader election, log replication, commit rules, and the failure scenarios that keep it honest.

Advanced · 22 min read

Why this matters

Every system you've met so far assumed someone is in charge — a leader to replicate from, a primary to fail over to. But who decides who's in charge when machines can crash, networks can split, and nobody has a global view? That problem is consensus, and it's the load-bearing wall under etcd, Consul, Kubernetes' control plane, and every serious distributed database's metadata layer. Raft (Ongaro & Ousterhout, 2014) is the consensus algorithm designed to be understandable — which, given that its predecessor Paxos has a reputation as "the algorithm nobody fully understands," is saying something.

Analogy: a committee electing a chairperson. The committee needs exactly one chair at a time — two chairs issuing contradictory orders is worse than none. Members vote, the chair sends regular "I'm still here" heartbeats, and if the heartbeats stop, someone calls for a new election. Raft is this committee with the bylaws written very, very carefully — because the failure mode isn't an awkward meeting, it's two leaders and a corrupted database.

Video: Consensus in Distributed Systems | Paxos & Raft — udaykiran․tech
Explains why distributed systems need consensus, the failures it handles, and the properties it must guarantee.

The three roles and the heartbeat

Every server is always in exactly one state: leader, follower, or candidate. Time is divided into terms (numbered eras, like "term 7"). The rules:

flowchart TD
    F[Follower<br/>term 7] -->|election timeout<br/>no heartbeat| C[Candidate<br/>term 8, votes self]
    C -->|majority of votes| L[Leader<br/>term 8]
    C -->|discovers higher term| F2[Follower<br/>term 9]
    C -->|timeout, split vote| C2[Candidate<br/>term 9, retry]
    L -->|higher term seen| F3[Follower]

The randomization of election timeouts is the quiet genius: if two followers time out simultaneously and both become candidates, they split the vote and both fail — then randomized timeouts make it overwhelmingly likely that one retries first and wins cleanly. Deterministic timeouts would livelock forever, like two people repeatedly saying "no, you go first" in a doorway.

Interactive diagram: StepThrough (loads in the app)

Video: Raft Consensus Explained: How Machines Agree Without a Boss — SysSketch
Covers leader, follower, and candidate roles, terms, and how randomized heartbeats keep the cluster stable.

Log replication: the actual work

Election just picks who's in charge. The point of consensus is agreeing on a log — an ordered sequence of commands every node applies identically. The leader accepts client writes, appends them to its log, and replicates via AppendEntries RPCs:

sequenceDiagram
    participant Client
    participant L as Leader (term 8)
    participant F1 as Follower 1
    participant F2 as Follower 2

    Client->>L: Write x=3
    L->>L: Append to log [index 42]
    par Replicate
        L->>F1: AppendEntries [42]: x=3
        L->>F2: AppendEntries [42]: x=3
    end
    F1-->>L: ack
    Note over L,F2: Majority (2 of 3) reached
    L->>L: COMMIT index 42
    L-->>Client: Success
    L->>F2: AppendEntries (next heartbeat carries commitIndex)

Interactive diagram: StepThrough (loads in the app)

The commit rule is everything: an entry is committed once a majority has it. A committed entry will survive any future election — Raft's election rules guarantee the new leader's log contains all committed entries (a candidate's log must be at least as up-to-date as a majority's, or it can't win). This is the safety property the whole algorithm exists to provide.

Video: Distributed Systems 6.2: Raft — Martin Kleppmann
Kleppmann's lecture walks through leader-driven AppendEntries log replication, consistency checks, and commit rules.

The numbers that matter

Back-of-envelope math for a typical 5-node etcd cluster:

Video: Distributed Systems 6.1: Consensus — Martin Kleppmann
Lays out quorum arithmetic, majority overlap, and how many failures a cluster can tolerate.

Failure scenarios (the exam questions)

Split brain that isn't: a network partition splits 5 nodes into 3+2. The side with 3 elects a leader and keeps committing; the side with 2 can't reach majority, so its would-be leader's writes never commit and clients get errors. No divergence, ever. The minority partition stops rather than going rogue. Compare with the AP systems from your CAP lesson and notice the price being paid: availability.

The slow follower: one follower's disk is dying, so its acks crawl. The leader only needs a majority — with 5 nodes, the two slowest nodes are simply not on the critical path. This is the deep reason odd-numbered quorums degrade gracefully: you always have slack for the stragglers.

Leader with a stale term: a partitioned old leader (term 8) reconnects and tries to send heartbeats. Everyone is on term 9; the higher term wins, and the old leader immediately steps down to follower. Terms are the algorithm's monotonic clock — they make "who's boss" unambiguous even across partitions and crashes.

Video: Split-Brain, Quorum, and Fencing | Leader Election — Ops and Odds
Walks through split-brain, dead leaders, zombie leaders, and fencing tokens that stop double writers.

When not to use it

Raft gives you linearizable consensus at the cost of write latency, a single-writer bottleneck, and operational care (quorum loss = full outage; back up your etcd). If you need coordination (locks, leader election, config) — use it, via etcd/Consul, don't hand-roll it. If you need throughput (product data at 100k writes/sec) — you want partitioned datastores with weaker consistency, not consensus. And never, ever implement Raft yourself for production unless implementing Raft is your product.

Video: "Consistency without consensus in production systems" by Peter Bourgon — Strange Loop Conference
Bourgon's Strange Loop talk shows consensus-free consistency with CRDTs and when to skip consensus entirely.

Takeaways

  1. Raft = leader election (majority vote, randomized timeouts) + log replication (commit on majority ack).
  2. 2F+1 nodes tolerate F failures; even node counts buy nothing — run 3 or 5.
  3. Commits cost one majority round trip (~RTT); elections cost ~one timeout (~0.5 s of write unavailability).
  4. The minority partition halts rather than diverging — that's the safety guarantee you're paying for.

Check your understanding

  1. In Raft, when is a log entry considered committed?

    • Once a majority of cluster nodes have replicated it
    • After the election timeout elapses
    • As soon as the leader appends it to its own log
    • Once all nodes have acknowledged it
  2. Why are Raft election timeouts randomized (150–300 ms) rather than fixed?

    • To make the algorithm harder for attackers to predict
    • Randomization is required for the log replication step
    • To reduce network bandwidth usage
    • To prevent split votes from livelocking when multiple candidates start elections simultaneously
  3. A 5-node Raft cluster splits 3/2 in a network partition. What happens?

    • Both sides elect leaders and continue accepting writes, diverging
    • The 3-node side keeps working; the 2-node side cannot reach majority and stops committing writes
    • The cluster shuts down entirely until the partition heals
    • The 2-node side takes over because it has less contention
  4. Why is running Raft with 4 nodes worse than 3?

    • 4 nodes cannot form a quorum at all
    • 4 nodes double the election timeout
    • Both tolerate 1 failure, but 4 needs 3 acks per commit instead of 2 — slower with no extra safety
    • Raft requires a prime number of nodes

Go deeper

Want to keep pulling this thread? These talks and tutorials go further than we did here:

Sources & further reading