SproutStackSproutStack home

Distributed Consensus: Paxos & Raft

3 min readBeginner-friendlyPremium
Premium chapter — free for early learners.

Read a little, play a little. No scary maths, and no rush.

Share:

The problem: who's right when nodes disagree?

Run a database on one machine and there's only ever one answer to "what's the current value?" Run it on five machines for durability, and a network hiccup can leave two machines thinking different things are true. Consensus is the protocol that lets a majority of nodes agree on one value — one leader, one committed log entry, one source of truth — even when some nodes are slow, crashed, or temporarily unreachable.

Why "majority," not "everyone"

Requiring all N nodes to agree means one dead machine halts the entire system — the opposite of the fault tolerance you built the cluster for. Requiring a quorum (a majority, e.g. 3 of 5) means the system keeps working as long as a majority is healthy, and — critically — any two quorums overlap by at least one node, which is what prevents two different values from both getting "accepted" during a split.

Paxos vs Raft, in one line each

Paxos proves consensus is possible and is famously hard to implement correctly from the paper alone — it's the theoretical foundation, not usually what you'd hand-roll. Raft was designed explicitly to be understandable: it elects a single leader (via randomized timeouts, so two nodes rarely race for it), the leader sequences all writes into a replicated log, and followers replay that log in order. Most consensus you'll touch in practice (etcd, Consul, CockroachDB's internals) is Raft or a close variant, specifically because it's easier to reason about and debug.

Where you actually meet this

You rarely implement consensus yourself — you depend on something that already has (etcd for Kubernetes' cluster state, ZooKeeper for older Kafka coordination, Raft inside CockroachDB or TiKV). What you need as an engineer is the mental model: a leader election happening somewhere means a brief window with no leader, writes during that window either queue or fail, and that's an availability/consistency tradeoff, not a bug.

Remember this

  • Consensus = majority agreement on one value, surviving node failures and delays.
  • Quorums overlap by design — that overlap is what prevents split-brain disagreement.
  • Raft won on understandability: single leader, replicated log, deliberately simple to reason about.

Check your understanding

2 questions · correct answers earn XP once each

1. Consensus algorithms exist to solve…
2. Why does a majority (quorum), not all nodes, need to agree?

My notes

Saved in this browser. Highlight a line above and save it, or write it in your own words.

Nothing saved yet. Your highlights will live here.

References

Finished reading?

Ticking it here also ticks the chapter in the sidebar, the section count and your streak — it is all one number.

Related chapters

Spotted a mistake or want a topic covered? Report an issue