Distributed Consensus (Raft & Paxos)
The Network Partition That Elected Two Leaders
Test your architecture intuition: Pitch a 7-axis solution, survive two aggressive reviewer objections, and inspect the staff-level Teacher Gold Answer.
1. What It Is & Why It Exists
The Core Problem: Agreement Across Unreliable Networks
In a distributed system, physical servers must coordinate state changes over an asynchronous, unreliable network subject to packet loss, arbitrary latency spikes, network partitions, and unannounced node crashes without Byzantine malice.
Coordinating state across such an environment without a single point of failure introduces the Distributed Consensus Problem: how can a cluster of independent machines agree on a single, deterministic sequence of values or state transitions?
The FLP Impossibility Theorem (Fischer, Lynch, Paterson, 1985)
The foundational theorem of distributed computing proves that in a purely asynchronous network, no deterministic consensus protocol can guarantee both Safety and Liveness in the presence of even a single unannounced fail-stop crash:
- Safety: Nothing bad happens (the cluster never commits two conflicting states or elects two active leaders).
- Liveness: Something good eventually happens (the cluster never deadlocks and always makes progress).
The Breakdown: Asynchronous Replication vs. 2PC vs. Quorum Consensus
Concrete Proof: The Split-Brain Blackout vs. Quorum Safety
Consider a 5-node cluster experiencing an asymmetric network partition dividing the cluster into two segments: a minority partition of 2 nodes () and a majority partition of 3 nodes ():
| Consensus Architecture | Behavior in Network Partition ( Split) | Safety Guarantee | Liveness Guarantee | Worst-Case Failure Mode |
|---|---|---|---|---|
| Primary-Backup (Async Replication) | Both partitions elect a primary. Both accept client writes. | 🚨 Zero Safety: State permanently diverges (Split-Brain). | ✅ High (Both sides accept writes) | Permanent data corruption requiring manual reconciliation. |
| Two-Phase Commit (2PC) | Coordinator cannot reach all 5 nodes; aborts all incoming transactions. | ✅ Safe (No split-brain) | 🚨 Zero Liveness: Cluster completely freezes if 1 node drops. | Total availability collapse under any single-node failure. |
| Quorum Consensus (Raft / Paxos) | Minority partition () can never commit a write (its clients time out); Majority partition () continues processing. | 🛡️ Linearizable Safety: At most 1 leader per term can secure a majority quorum. | 🛡️ Continuous Liveness: Survives node failures. | Minority partition blocks writes until network heals. |
The First-Principles Solution: Quorum-Based Replicated State Machines
Practical consensus protocols (Raft, Multi-Paxos, Zab) circumvent the FLP Impossibility Theorem by introducing partial synchrony (randomized timers and heartbeat timeouts):
- Absolute Safety: Safety is mathematically guaranteed under all asynchronous conditions, regardless of network delays or packet reordering.
- Conditional Liveness: Liveness is guaranteed as long as a strict mathematical majority quorum of nodes () can communicate, with messages arriving well within the election timeout often enough for a leader to be elected and stay elected.
- Replicated State Machine (RSM) Model: If identical deterministic state machines on multiple servers process an identical, ordered sequence of inputs from a consensus log, they will produce an identical, consistent state across all nodes.
Synthesizing vector architecture diagram...
Both panels start with a five-node cluster split by a network partition into groups of 2 and 3. In the "Broken Baseline" panel, there is no quorum rule, so each side elects its own leader, Node 1 accepts x=10 and Node 3 accepts x=20, and when the network heals the two histories conflict and one side's writes must be thrown away. In the "Production Standard" panel, Raft requires 3 of 5 votes: the 2-node minority's old leader cannot get a quorum, so it may keep appending writes to its log but can never commit them (its clients time out, and those entries are deleted when it rejoins), while the 3-node majority elects Node 3 as leader for term 2 and keeps committing. A cluster of 2F+1 nodes survives F failures, and a minority can never commit, so the two sides can never disagree.
2. Core Mechanics & Algorithmic Architecture
The Raft Algorithm: 3 Decomposed Sub-Mechanisms
Raft decomposes the consensus problem into three strictly formalized sub-problems: Leader Election, Log Replication, and Safety.
Synthesizing vector architecture diagram...
Every node boots as a Follower and waits for heartbeats. If none arrives within its randomized election timeout (150-300 ms in the Raft paper's example; etcd's default is 1,000 ms, randomized between 1x and 2x), it becomes a Candidate, increments its term, votes for itself, writes its term and vote to disk, and asks the others for votes. From Candidate there are three ways out: a majority of votes makes it Leader; learning of a valid leader with a current or higher term sends it back to Follower; and a split vote makes it time out and start again with a higher term. A Leader stays leader until it sees any message with a higher term, then immediately steps down to Follower. The term number is Raft's logical clock: a higher term always wins, which is what guarantees at most one leader per term.
1. Leader Election & Mathematical Quorum
- Nodes start as Followers. If a follower receives no heartbeat within a randomized
electionTimeout( in the Raft paper's example; etcd defaults to , randomized up to twice that), it transitions to Candidate, increments itscurrentTerm, votes for itself, writescurrentTermandvotedForto disk (otherwise a restarted node could vote twice in one term), and broadcastsRequestVoteRPCs. - Election restriction: a voter refuses a candidate whose log is less up-to-date than its own (a later last term wins; with equal last terms, the longer log wins). A committed entry is on a majority and a winner needs a majority of votes, so the two share a node, and that node refuses any candidate missing the entry: every committed entry survives an election.
- A Candidate becomes Leader only after securing votes from a strict mathematical quorum:
- For , (tolerates node failure).
- For , (tolerates node failures).
- An even-numbered cluster () requires votes, providing the exact same fault tolerance () as while increasing network overhead. Always deploy odd-numbered clusters ().
2. Log Replication & The Log Matching Property
- The Leader receives commands from clients, appends them to its local Write-Ahead Log, and broadcasts
AppendEntriesRPCs. - An entry from the leader's current term is Committed once it is safely written on a majority of nodes (). Raft never commits entries from earlier terms by counting replicas: a new leader first appends and commits a no-op entry in its own term, and every earlier entry becomes committed with it. Until that no-op commits, a new leader doesn't know the commit index and can't serve linearizable reads (the section The leader crashes: elections of Replication, Quorums & Read-Your-Writes).
- Log Matching Invariant: If two distinct logs contain an entry with the same index and term, they are guaranteed to store identical commands in all entries up through that index.
3. Linearizable Reads (Avoiding Stale Reads on Partitioned Leaders)
If a Leader is partitioned off, it might respond to read queries with stale data before discovering its isolation. Raft guarantees linearizable reads via two protocols:
- ReadIndex Protocol: The Leader records its current
commitIndex, broadcasts a heartbeat to confirm it still commands a majority quorum, and serves the read once ACKs arrive and it has applied up to the recorded index. A deposed leader cannot collect a majority, so its read fails instead of returning stale data. - LeaseRead Protocol: The Leader relies on a physical time-bounded lease, timed from the send time of the last heartbeat a majority acknowledged: as long as , it serves reads locally with zero network hops. must be shorter than the minimum election timeout minus a margin for the difference in clock rates, so the lease ends before any other node can win an election. Unlike ReadIndex, this relies on clocks and bounded pauses for safety (the section A network partition and two leaders of Replication, Quorums & Read-Your-Writes; clock margins in the section Time: whose clock decides? of Leases, Fencing Tokens & Distributed Locks).
Step-by-Step Deterministic Write & Replication Pipeline
Synthesizing vector architecture diagram...
In-Memory & Distributed State Structures
Every Raft node maintains three distinct classes of state in memory and on disk:
| State Category | Variable Name | Type / Storage Medium | Purpose |
|---|---|---|---|
| Persistent on All Nodes | currentTerm | int64 (WAL on NVMe) | Latest term server has seen; monotonically increases |
| Persistent on All Nodes | votedFor | node_id (WAL on NVMe) | Candidate that received vote in current term (null if none) |
| Persistent on All Nodes | log[] | Array of entries (WAL) | Log entries containing command, index, and term |
| Volatile on All Nodes | commitIndex | int64 (RAM) | Index of highest log entry known to be committed |
| Volatile on All Nodes | lastApplied | int64 (RAM) | Index of highest log entry applied to state machine |
| Volatile on Leader Only | nextIndex[] | Array of int64 (RAM) | For each peer, index of next log entry to send |
| Volatile on Leader Only | matchIndex[] | Array of int64 (RAM) | For each peer, index of highest log entry known replicated |
Consensus Scenario & Quorum State Matrix
Under a 5-node cluster () with quorum requirement , and initial Term 2 Leader :
| Step # | Event / Input | In-Memory / Distributed State | Evaluation & Transition | Outcome / Cluster State |
|---|---|---|---|---|
| 1 | Client write to Leader SET balance = 500 | 5 nodes active, Term 2; Quorum threshold | writes Index 10, broadcasts AppendEntries;reply ACK ( slow, down). Total ACKs: 3 | Commit Succeeded entry committed; applies to state machine and returns 200 OK |
| 2 | Asymmetric network partition isolates from | Minority (2 nodes); Majority (3 nodes) | Client sends write to ; can only get 2 ACKs () blocked. election timeout trips (180ms); starts Term 3 election | Split-Brain Prevented can append but never commit (its clients time out); secures 3/3 votes and becomes legitimate Term 3 Leader |
| 3 | Network partition heals; all 5 nodes re-establish TCP | at Term 2 (stale uncommitted entry); at Term 3 (authoritative leader) | sends heartbeat to ; rejects with higher term (Term 3 > Term 2).steps down to Follower; syncs authoritative log | Authoritative Log Overwrite and overwrite uncommitted Term 2 entries with Term 3; cluster unified |
3. Data Migration, Anti-Entropy & Consistency Protocols
1. Log Compaction & Snapshotting
Because an append-only log grows indefinitely, servers would eventually run out of disk space and take hours to replay logs upon reboot. Raft handles this via Memory Snapshotting:
- An application periodically takes a point-in-time snapshot of its state machine (e.g., at Index 10,000).
- All log entries up through Index 10,000 are deleted from disk.
- The
InstallSnapshotRPC: If a lagging follower or newly added node is so far behind that its missing entries have already been compacted, the Leader streams its latest snapshot file directly to the follower over the network.
Synthesizing vector architecture diagram...
2. Dynamic Cluster Membership Changes (Joint Consensus)
Adding or removing nodes cannot be done via static configuration changes on all machines, as two independent majorities could vote at the same time during the rollout window.
- Raft Joint Consensus: The cluster transitions through an intermediate phase () where decisions require independent majorities from both the old configuration and the new configuration:
- Once is committed, the leader commits the final configuration .
3. Accelerated Log Conflict Resolution & Graceful Draining
- Accelerated Log Backtracking: In naive Raft, if a follower's log diverges across 1,000 entries, the leader decrements
nextIndexby 1 per rejectedAppendEntriesRPC (costing 1,000 network round-trips). In the optimized protocol, the follower returns the conflicting entry's term (conflictTerm) and the earliest index of that term (conflictIndex). The leader skipsnextIndexpast the entire conflicting term in a single RPC, accelerating log synchronization after network partition healing. - Graceful Leader Step-Down & Node Draining: When an active leader node is selected for host maintenance or rolling restarts:
- Operator triggers
etcdctl move-leader <target_follower_id>(or issues RaftTimeoutNowRPC). - The leader halts client writes, synchronizes uncommitted log entries to the target follower, and commands the follower to initiate an immediate election without waiting for election timeout.
- The target follower wins the election at once, so writes pause only briefly instead of for a full election timeout.
- The retiring node drains client gRPC connections and shuts down gracefully without triggering cluster-wide election storms.
- Operator triggers
Unlock Complete Architecture & Production Runbooks
You have explored the free architectural preview (~46%). Spend 1 Coin to unlock the remaining 7 production deep-dive sections for a full 24 hours.