Distributed Consensus
Distributed Consensus
Definition: Protocols that allow a collection of independent nodes in a distributed system to agree on a single data value or sequence of state machine commands, even in the presence of node failures or network delays.
How It Works
- Leader election: the cluster elects one node as leader, responsible for ordering and proposing new entries to a replicated log. Elections happen at startup and again any time the current leader is suspected dead (via missed heartbeats).
- Log replication: the leader accepts client writes, appends each as an entry to its local log, and replicates that entry to followers. An entry is committed once a majority (quorum) of nodes have durably stored it.
- Quorum requirement: with N nodes, a majority is
floor(N/2) + 1. Requiring a majority (not all nodes) means the system tolerates up tofloor((N-1)/2)node failures while still making progress, since any two majorities out of N nodes must overlap by at least one node. - Term/epoch numbers: each election increments a monotonic counter (Raft calls it a “term,” Paxos a “ballot number” or “proposal number”). This lets nodes detect and reject messages from a stale, previously-deposed leader, which is the core defense against split-brain.
- Raft: decomposes consensus into leader election, log replication, and safety rules, explicitly designed to be more understandable than Paxos. A candidate needs votes from a majority to become leader; only logs at least as up to date as a majority’s can win an election, which guarantees a new leader already has every committed entry.
- Paxos / Multi-Paxos: the original consensus proof, structured as proposers, acceptors, and learners running a two-phase (prepare/promise, then accept/accepted) protocol per log entry. Multi-Paxos optimizes repeated Paxos rounds by keeping a stable leader instead of re-running the prepare phase for every entry.
- ZAB (ZooKeeper Atomic Broadcast): ZooKeeper’s own consensus protocol, similar in spirit to Multi-Paxos with a strong emphasis on primary-order broadcast for its specific coordination use case.
- Once a majority has replicated an entry, the leader considers it committed, applies it to its own state machine, and informs followers to apply it too. A client only gets a success response after commit, guaranteeing the write survives any minority-side failure.
- Safety versus liveness: consensus protocols are designed to guarantee safety unconditionally (the system never agrees on two different values for the same log position, even during partitions or arbitrary delays) while only guaranteeing liveness (the system eventually makes progress) under weaker assumptions, typically that the network eventually stabilizes long enough for a leader to be elected and replicate an entry.
- This safety-first design is deliberate: a consensus protocol that occasionally sacrifices liveness (the cluster briefly can’t commit anything) is recoverable once conditions improve; one that sacrifices safety (two nodes disagree about committed history) can produce silent, unrecoverable data corruption.
- Fencing tokens: a monotonically increasing number issued alongside leadership, used by downstream systems to reject writes from a leader that doesn’t realize it’s been deposed yet, closing a race condition that term numbers alone don’t fully cover for external resources.
- Worked example, Raft term and log index mechanics: a 5-node cluster starts in term 1 with node A as leader. A writes log entries at indexes 1, 2, 3, each tagged
(term=1, index=N), and replicates them to a majority before committing. A crashes. The remaining nodes time out waiting for a heartbeat, increment their local term to 2, and hold an election: node B requests votes, presenting its log’s last entry as(term=1, index=3). Because B’s log is at least as up to date as a majority’s, it wins and becomes leader for term 2. - Node A recovers and rejoins as a follower. It had a fourth entry
(term=1, index=4)that it had appended locally but never replicated to a majority before crashing, so it was never committed. Raft’s log-matching rule forces A to discard that uncommitted entry and adopt B’s log instead, which is why only committed entries are ever guaranteed durable, an appended-but-unreplicated entry can be silently rolled back. - B now accepts new client writes as
(term=2, index=4),(term=2, index=5), and so on. Any message A might still send using its stale term-1 leadership claim is rejected by every other node, since they’ve all moved to term 2, this term comparison is the exact mechanism that prevents the old leader from causing a split-brain write.
Worked Example: Quorum Math at Different Cluster Sizes
| Cluster size (N) | Majority needed | Failures tolerated | Notes |
|---|---|---|---|
| 3 | 2 | 1 | Smallest practical fault-tolerant cluster |
| 4 | 3 | 1 | Same fault tolerance as N=3, costs one extra node and vote |
| 5 | 3 | 2 | Common production default, balances tolerance and quorum latency |
| 7 | 4 | 3 | Higher tolerance, but every write needs 4 responses instead of 3 |
Adding a node from an odd size to the next even size (3→4, 5→6) never improves fault tolerance, it only adds cost and a larger quorum to wait for, which is why production consensus clusters almost always stick to odd sizes.
Trade-offs
- Consensus provides strong consistency and fault tolerance for the specific data it manages, at the cost of requiring a majority round trip for every write, which caps throughput and adds latency compared to a single, unreplicated node.
- More nodes in the consensus group increase fault tolerance (more failures survivable) but also increase the size of the majority needed and the network cost per write, so consensus clusters are typically small (3, 5, or 7 nodes), not scaled out like a sharded data store.
- Consensus solves agreement on a single, small, frequently-updated log; it deliberately isn’t used to store bulk application data, because every byte written goes through majority replication and leader serialization, which doesn’t scale the way sharding does.
- Raft trades some of Paxos’s generality and flexibility for understandability and a directly implementable specification, which is why it displaced Paxos in most new systems despite Paxos being proven first and being marginally more flexible in exotic configurations.
Why It Matters
- It’s the mechanism underneath distributed coordination primitives: leader election, distributed locks, configuration management, and service discovery all reduce to “get a cluster of nodes to agree on one value.”
- It’s what makes a CP data store actually CP: consensus is the algorithm that guarantees the minority side can’t diverge from the majority, which is the concrete implementation of “reject requests rather than serve stale data.”
- It underlies the control plane of most modern infrastructure: Kubernetes’ cluster state, service configuration, and leader election for many distributed databases all depend on a consensus protocol running underneath.
Common Pitfalls
- Split-brain: a network partition leads to two sub-clusters, each believing it’s the leader, if quorum rules are violated (e.g., a badly configured system that allows a minority partition to keep accepting writes).
- Confusing “commit” with “apply.” An entry is committed once a majority has it durably stored; it’s applied to the state machine (and visible to clients) potentially slightly after that. Reading directly from a follower without checking commit status can return an uncommitted, later-reverted entry.
- Running an even number of nodes. An even-sized cluster doesn’t add fault tolerance over the next smaller odd number (5 nodes tolerates 2 failures, same as needing a majority of 6 which also tolerates only 2), it just adds cost and a higher chance of split votes during elections.
- Treating consensus as a general-purpose database. It’s built for small, frequently-agreed-upon state (config, leader identity, locks), not for storing gigabytes of application data, because every write pays the full majority round-trip cost.
- Assuming a leader is dead just because it’s slow. Aggressive election timeouts can cause unnecessary leadership churn (“flapping”) under normal network jitter or GC pauses, which itself hurts availability.
- Deploying a consensus cluster’s members across an even split of failure domains (e.g., exactly 2 nodes in each of two racks/AZs plus 1 elsewhere) without checking that no single failure domain’s loss can strand the rest below a majority.
- Ignoring disk fsync guarantees. If a node acknowledges a log write without actually forcing it to durable storage, a power loss can silently violate the “committed entries survive crashes” guarantee the whole protocol depends on.
Comparison
| Raft | Paxos / Multi-Paxos | ZAB (ZooKeeper) | |
|---|---|---|---|
| Primary goal | Understandable, implementable consensus | Original proven consensus protocol | Primary-order atomic broadcast |
| Leader model | Explicit, strong leader | Implicit via stable-leader optimization | Explicit, strong leader |
| Common implementations | etcd, Consul, CockroachDB, TiKV | Chubby, Spanner (Paxos variants) | ZooKeeper |
| Reputation | Easier to reason about and implement correctly | Notoriously hard to implement correctly from the paper alone | Purpose-built for ZooKeeper’s needs |
| Log structure | Explicit replicated log with term/index per entry | Per-slot decisions, log built from separate agreements | Sequential zxid-tagged transactions |
| Handles configuration changes (membership) | Joint consensus, defined in the core spec | Requires separate reconfiguration protocols, historically tricky | Supported via its own reconfiguration extension |
Real-World Scenario
A Kubernetes cluster’s control plane state (which pods are scheduled where, current deployment state) lives in etcd, a Raft-based key-value store. When the etcd leader node crashes, the remaining nodes detect missed heartbeats, time out, and hold a new election. A follower with an up-to-date log wins a majority vote and becomes the new leader within a second or two. During that brief election window, the Kubernetes API server’s writes to etcd are rejected or queued, but no committed cluster state is lost, because every previously committed entry was already replicated to a majority, including at least one node in the new leader’s majority.
Debugging Walkthrough: A Failed Election and Near Split-Brain
- A 5-node Raft cluster spans two racks: 3 nodes on rack A, 2 on rack B. A top-of-rack switch issue partially isolates rack B, delaying (not dropping) packets between racks by several seconds.
- The leader happens to be on rack B. Nodes on rack A stop receiving its heartbeats within their election timeout, increment their term, and one of them (node A1) starts an election. It gets votes from the other two rack-A nodes, a majority (3 of 5), and becomes leader for the new term.
- Meanwhile, the original leader on rack B hasn’t realized anything is wrong yet, its own clock hasn’t hit its own timeout, and it keeps sending heartbeats and trying to replicate entries, but its RPCs to rack-A nodes now get rejected because those nodes are on a higher term.
- On-call sees two things simultaneously: a leader-election metric firing (new leader elected on rack A), and the old leader’s logs showing a burst of “stale term, rejecting append” errors. This is not yet split-brain, just a leadership handoff that the old leader hasn’t caught up to.
- The critical check: did the old leader ever manage to commit a write during the confusion? Because committing requires a majority, and it can no longer reach a majority once rack A moved to a new term, it can’t have committed anything new after the network issue started, so client writes accepted during the ambiguous window either succeeded against the new leader or were correctly rejected by the old one.
- Once the network delay clears, the old leader receives an append or heartbeat carrying the newer term number, recognizes it’s stale, and steps down to follower, catching up on any entries it’s missing.
- The postmortem confirms this was consensus working as designed, brief unavailability for writes on the minority side, no data loss or divergence, rather than an actual split-brain, which would only happen if a system violated the majority rule (e.g., misconfigured to allow a minority to also elect a leader).
FAQ
Is distributed consensus the same as distributed transactions? No. Consensus is about a cluster of replicas of the same logical value agreeing on its next state; distributed transactions are about coordinating an atomic operation across different pieces of data, often on different services or shards.
Why do consensus clusters use odd numbers of nodes? Because an odd count gives the same fault tolerance as the next-larger even count at lower cost; 3 and 4 nodes both tolerate only 1 failure, but 4 requires one more node’s cost and vote.
Can consensus systems scale writes horizontally? Not directly. Since every write must go through a single leader and a majority quorum, adding more nodes to the consensus group increases fault tolerance, not write throughput; scaling write throughput requires sharding into multiple independent consensus groups.
Can a consensus system lose data that was already acknowledged to a client? Not if implemented correctly, an entry is only acknowledged after majority durability, and any future leader is guaranteed (by the election rule) to already have it. Data loss would indicate a bug or a violated assumption, like disks not actually being durable.
What happens if the network partitions a consensus cluster exactly in half? Neither half has a majority, so neither can elect a leader or commit writes, and the whole system becomes unavailable for writes until the partition heals or enough nodes rejoin one side.
Common Interview Questions
- Why does consensus require a majority rather than all nodes? Requiring all nodes would mean a single failure halts the system; a majority tolerates minority failures while guaranteeing any two majorities overlap, preserving agreement.
- What’s the practical difference between Paxos and Raft? They solve the same problem, but Raft explicitly separates leader election, log replication, and safety into named sub-protocols to make correct implementation more tractable; Paxos’s single-decree formulation is more general but harder to extend correctly to a replicated log.
- How does a new leader know it has all committed entries? Election rules require a candidate’s log to be at least as up to date as a majority of voters’, which guarantees any log entry committed by a prior majority is present in the new leader’s log too.
- What happens to in-flight writes during a leader election? They’re rejected or time out; clients must retry against the new leader once one is elected, since the old leader can no longer safely commit anything.
- Why is consensus considered expensive compared to plain asynchronous replication? Every write requires a majority round trip before being acknowledged, versus asynchronous replication’s single-node write followed by best-effort propagation.
Example
etcd (used by Kubernetes for cluster state) and HashiCorp Consul both use the Raft consensus algorithm. Google’s Chubby lock service and Spanner use Paxos-family protocols. Apache ZooKeeper uses ZAB, its own purpose-built consensus protocol, for coordination primitives like leader election and distributed locks.
Design Checklist
- Does this component actually need linearizable agreement, or would eventually-consistent replication be enough? Consensus is expensive; reach for it only when correctness genuinely requires it (leader election, config, locks).
- Is the consensus group sized correctly (odd number, small enough to keep quorum latency low)?
- What’s the plan for a network partition that splits the cluster with no majority on either side? The honest answer should be “the system becomes unavailable for writes until it heals,” not a workaround that risks split-brain.
- Are downstream systems that trust “who’s the leader” protected against a stale, deposed leader via fencing tokens or equivalent?
Related Terms
Referenced by