Learn Labs
10. Consistency and Consensus

10.3 Consensus

The three things that are easy on one node and hard with fault tolerance:

  1. Linearizable database — easy with one leader; but how do you FAIL OVER while avoiding split brain? How do you ensure a node that believes it's the leader hasn't been voted out while temporarily paused?
  2. Linearizable ID generator — just a counter with atomic fetch-and-add; what if it crashes?
  3. Atomic CAS — may be one CPU instruction on a single node; how do you make it fault-tolerant?

It turns out ALL OF THESE ARE INSTANCES OF THE SAME FUNDAMENTAL PROBLEM: CONSENSUS.

The four best-known algorithms: Viewstamped Replication, Paxos, Raft, Zab. "These algorithms have quite a few similarities, but THEY ARE NOT THE SAME." All work in a non-Byzantine system model — network communication may be arbitrarily delayed or dropped, nodes may crash/restart/disconnect, but nodes otherwise follow the protocol and do not behave maliciously.

3.1 The FLP result — what it actually says

The FLP result proves NO ALGORITHM IS ALWAYS ABLE TO REACH CONSENSUS if there is a risk that a node may crash. Yet here we are, discussing consensus algorithms. What's going on?

First, FLP DOESN'T SAY WE CAN NEVER REACH CONSENSUS; it only says WE CAN'T GUARANTEE A CONSENSUS ALGORITHM WILL ALWAYS TERMINATE.

Second, FLP is proved assuming a DETERMINISTIC algorithm in the ASYNCHRONOUS system model — meaning THE ALGORITHM CANNOT USE ANY CLOCKS OR TIMEOUTS. If it CAN use timeouts to suspect a node has crashed (even if the suspicion is sometimes wrong), CONSENSUS BECOMES SOLVABLE. Even allowing RANDOM NUMBERS is sufficient.

3.2 The many faces of consensus — all equivalent

All of these are equivalent. If you have an algorithm solving one, you can convert it into a solution for any of the others.

  • Single-value consensus ⟺ linearizable CAS
  • Shared log — total order broadcast, also called atomic broadcast
  • Atomic commitment — 2PC's problem
  • Atomic fetch-and-add — only almost equivalent: it works for two nodes only, so its consensus number is 2
(a) Single-value consensus — the four properties
PropertyDefinitionKind
Uniform agreementNo two nodes decide differentlySAFETY
IntegrityAfter a node has decided one value, it cannot change its mindSAFETY
ValidityIf a node decides value v, then v was PROPOSED BY A NODESAFETY
TerminationEvery node that does not crash EVENTUALLY DECIDES a valueLIVENESS

Agreement and integrity define the core idea. VALIDITY RULES OUT TRIVIAL SOLUTIONS — an algorithm that always decides null would satisfy agreement and integrity, but not validity.

If you don't care about fault tolerance, the first three are EASY: hardcode one node as the "DICTATOR." But if that node fails, the system can no longer decide anything. ALL THE DIFFICULTY ARISES FROM THE NEED FOR FAULT TOLERANCE.

TERMINATION formalizes fault tolerance: the algorithm must MAKE PROGRESS even if some nodes fail. And it must decide EVEN IF A CRASHED NODE NEVER COMES BACK. (Instead of a software crash, imagine an earthquake causes the datacenter to be destroyed by a landslide — you must assume your node is buried under 30 feet of mud and is never coming back online.)

Any consensus algorithm requires AT LEAST A MAJORITY OF NODES to be functioning correctly in order to assure termination.

However, most consensus algorithms ensure THE SAFETY PROPERTIES ARE ALWAYS MET, EVEN IF A MAJORITY OF NODES FAIL or a severe network problem occurs. Thus a large-scale outage can STOP THE SYSTEM FROM PROCESSING REQUESTS, BUT IT CANNOT CORRUPT THE CONSENSUS SYSTEM BY CAUSING INCONSISTENT DECISIONS.

(b) CAS ⟺ consensus
Cas ⇒ consensus
set the object to NULL. Each node CASes with expected=null and new=its proposed value. The decided value is whatever the object ends up set to.
Consensus ⇒ cas
when nodes want to CAS with the same expected value, use consensus to propose the new values, then set the object to whatever was decided. CAS invocations whose value wasn’t decided Return an error. Different expected values use Separate runs of the protocol.

(A real example: conditional writes in object stores — Ch 6 §1.3.)

(c) Shared logs ⟺ consensus

The five properties of a shared log:

PropertyDefinition
Eventual appendIf a node requests a value be added and doesn't crash, it must EVENTUALLY READ that value in a log entry
Reliable deliveryNo entries are lost — if one node reads an entry, eventually EVERY non-crashed node reads it
Append-onlyAfter a node reads an entry it is IMMUTABLE, and new entries can be added only AFTER it, not before. Rereading gives the same entries in the same order, even after a crash and restart
AgreementIf two nodes both read entry e, then PRIOR TO e THEY MUST HAVE READ EXACTLY THE SAME SEQUENCE of entries in the same order
ValidityIf a node reads an entry containing a value, a node PREVIOUSLY REQUESTED that value's addition

(A shared log is implemented with a TOTAL ORDER BROADCAST protocol, also known as atomic broadcast or total order multicast: to add a value you "broadcast" it; when the protocol "delivers" it, it becomes a log entry.)

Shared log ⇒ consensus
every node requests its value be added; Whichever value is read back in the first log entry is decided. Since all nodes read entries in the same order, they agree.
Consensus ⇒ shared log
  1. a Slot for every future entry; run a Separate instance of consensus per slot
  2. to add a value, propose it for an Undecided slot
  3. when a slot is decided And all previous slots are decided, append it — plus any consecutive decided slots
  4. if your value wasn’t chosen, Retry with a Later slot

Single-leader replication WITHOUT failover does not meet the liveness requirement, since it stops delivering messages if the leader crashes. AS USUAL, THE CHALLENGE IS PERFORMING FAILOVER SAFELY AND AUTOMATICALLY.

(d) Fetch-and-add — the one that ALMOST works
Cas ⇒ fetch-and-add
read, then CAS(expected=what you read, new=that+1); retry on failure. Less efficient under contention, Functionally equivalent.
Fetch-and-add ⇒ consensus?

Initialize to 0; every proposer increments.

  • One node reads 0 — call it the winner, its value is decided.
  • But the others know they are Not the winner and Don’t know which node won.
  • The winner could tell them — But what if the winner crashes first?
  • The others are Left hanging, unable to decide. And they Can’t fall back to another node, because The node that read 0 may yet come back and rightly decide its value.

⇒ Termination fails.

Exception: with At most two proposers, they exchange values first; the node reading 0 decides its own, the node reading 1 decides the other’s.

⇒ Fetch-and-add has a consensus number of 2. Cas and shared logs have a consensus number of ∞.

(e) Atomic commitment ⟺ consensus

The important difference: WITH CONSENSUS IT'S OK TO DECIDE ANY VALUE THAT WAS PROPOSED, WHEREAS WITH ATOMIC COMMITMENT THE ALGORITHM MUST ABORT IF ANY PARTICIPANT VOTED TO ABORT.

Atomic commitment's five properties: uniform agreement · integrity · validity — "If a node commits, ALL nodes must have previously voted to commit. If ANY node voted to abort, ALL nodes must abort" · nontriviality — "If all nodes vote to commit, AND NO COMMUNICATION TIMEOUTS OCCUR, then all nodes must commit" (this rules out an algorithm that always aborts) · termination.

Consensus ⇒ atomic commit: every node sends its vote to every other node. Nodes receiving commit-votes from Everyone propose “commit” via consensus; nodes receiving an abort-vote Or experiencing a Timeout propose “abort.” Then commit or abort per the consensus decision.

⇒ “commit” is proposed Only if all nodes voted to commit. It could happen that some propose abort and others commit if all voted commit but Some communication timed out — In which case it doesn’t matter which, as long as they all do the same thing.

3.3 Consensus in practice — shared logs win

Which formulation is most useful in practice? MOST CONSENSUS SYSTEMS PROVIDE SHARED LOGS. Raft, Viewstamped Replication, and Zab provide them out of the box. Paxos provides single-value consensus, but in practice most systems using Paxos use MULTI-PAXOS, which also provides a shared log.

Why a shared log is such a good primitive:

UseHow
Database replicationEvery entry is a write; every replica processes the same writes in the same order with DETERMINISTIC logic ⇒ all replicas end up consistent. This is STATE MACHINE REPLICATION, the principle behind event sourcing (Ch 3)
Serializable transactionsEvery entry is a DETERMINISTIC transaction executed as a stored procedure; every node executes them in the same order ⇒ SERIALIZABLE (Ch 8 §4.1)
Single-value consensus / CASDecide the value that appears FIRST in the log
Many instances of single-value consensus (one per theater seat)Include the seat number in the entry; decide the FIRST entry containing that seat number
Atomic fetch-and-addPut the addend in an entry; the current value is the SUM of all entries so far
Fencing tokensA simple counter on log entries. In ZooKeeper this is the zxid
Stream processingCh 12

⚠️ Sharded databases with a strong consistency model often maintain A SEPARATE LOG PER SHARD, which improves scalability BUT LIMITS THE CONSISTENCY GUARANTEES (consistent snapshots, foreign-key references) they can offer ACROSS shards.

3.4 From single-leader replication to consensus — breaking the circularity

It seems we need consensus to elect a leader, and we need a leader to solve consensus. HOW DO WE BREAK OUT OF THIS CONUNDRUM?

Consensus algorithms DON'T REQUIRE THAT THERE IS ONLY ONE LEADER AT ANY ONE TIME. Instead they make a WEAKER guarantee: they define an EPOCH NUMBER and guarantee that WITHIN EACH EPOCH, THE LEADER IS UNIQUE.

The epoch number under four names:

AlgorithmName
Paxosballot number
Viewstamped Replicationview number
Raftterm number
(generic)epoch number

The two-round structure:

  1. Round 1 — elect a leader. A node that hasn't heard from the leader for some timeout starts a vote. The election gets a new epoch number greater than any previous one. “If a conflict arises between two leaders in two epochs (perhaps because the previous leader wasn't dead after all), the leader with the higher epoch number prevails.”
  2. Round 2 — vote on each log entry. Before appending, the leader must check there isn't another leader with a higher epoch. It collects votes from a quorum. “A node votes yes only if it is not aware of any other leader with a higher epoch.”

★ The quorums for the two votes must overlap: if a proposal vote succeeds, at least one of the nodes that voted for it must also have participated in the most recent successful leader election. ⇒ If the proposal passes without revealing a higher-numbered epoch, the leader can conclude no higher-epoch leader has been elected, and can safely append.

These two rounds look SUPERFICIALLY SIMILAR TO 2PC, BUT THEY ARE VERY DIFFERENT PROTOCOLS:

Consensus2PC
Who starts itANY node can start an electionONLY the coordinator can request votes
What's requiredOnly a QUORUM of nodes to respondA YES VOTE FROM EVERY PARTICIPANT

3.5 Subtleties

Every new log entry is synchronously replicated to a quorum before it is confirmed to the client — this ensures the entry won't be lost if the current leader fails.

How the algorithms differ on honoring the old leader's entries:

AlgorithmApproach
RaftAllows a node to become leader ONLY IF ITS LOG IS AT LEAST AS UP TO DATE as those of a majority of its followers
PaxosAllows ANY node to become leader, but REQUIRES IT TO BRING ITS LOG UP TO DATE with other nodes before appending new entries

⚠️ The consistency-vs-availability choice inside leader election:

It's ESSENTIAL that the new leader is up to date with any confirmed entries before processing writes or linearizable reads. IF A NODE WITH STALE DATA BECAME LEADER, IT MIGHT WRITE NEW VALUES TO LOG ENTRIES THAT WERE ALREADY WRITTEN by the old leader, VIOLATING THE APPEND-ONLY PROPERTY.

In some cases you might choose to WEAKEN the consensus properties to recover more quickly, or to be able to recover at all. Kafka offers UNCLEAN LEADER ELECTION, allowing ANY replica to become leader even if not up to date. Also, in databases with asynchronous replication, you cannot guarantee ANY follower is up to date when the leader fails.

If you drop the requirement, you may improve performance and availability, BUT YOU ARE ON THIN ICE, SINCE THE THEORY OF CONSENSUS NO LONGER APPLIES. While things will work fine as long as there are no faults, the problems in Ch 9 CAN EASILY CAUSE DATA LOSS OR CORRUPTION.

Linearizable reads need a quorum too:

Turning writes into log entries and replicating them to a quorum ISN'T ALL THAT'S REQUIRED. If you want LINEARIZABLE READS, THEY ALSO HAVE TO GO THROUGH A QUORUM VOTE, similarly to a write, TO CONFIRM THAT THE NODE THAT BELIEVES ITSELF TO BE LEADER REALLY IS STILL UP TO DATE. Linearizable reads in etcd work like this.

Reconfiguration: most algorithms in standard form assume a FIXED SET OF NODES. Extensions make adding/removing nodes possible — especially useful when adding new regions, or MIGRATING from one location to another (first adding new nodes, then removing old ones).

3.6 Pros and cons

Consensus is essentially "SINGLE-LEADER REPLICATION DONE RIGHT," with automatic failover on leader failure, ensuring that NO COMMITTED DATA IS LOST and SPLIT BRAIN IS NOT POSSIBLE, even in the face of all the problems in Ch 9.

ANY SYSTEM THAT PROVIDES AUTOMATIC FAILOVER BUT DOES NOT USE A PROVEN CONSENSUS ALGORITHM IS LIKELY TO BE UNSAFE. (Using one is not a guarantee of whole-system correctness — there are still plenty of places where bugs can lurk — but it's a good start.)

The five costs:

  1. Always requires a STRICT MAJORITY — three nodes to tolerate one failure, five to tolerate two
  2. Every operation requires communication with a quorum, so YOU CAN'T INCREASE THROUGHPUT BY ADDING MORE NODES — IN FACT, EVERY NODE YOU ADD MAKES THE ALGORITHM SLOWER
  3. If a network partition cuts off some nodes, ONLY THE MAJORITY PORTION CAN MAKE PROGRESS; the other nodes are BLOCKED
  4. Timeout tuning is hard in environments with highly variable network delays, especially across regions. Too large → slow recovery. Too small → lots of unnecessary leader elections, resulting in TERRIBLE PERFORMANCE as the system SPENDS MORE TIME CHOOSING LEADERS THAN DOING USEFUL WORK.
  5. Sensitivity to specific network problems. Raft has unpleasant edge cases: if the entire network works correctly EXCEPT ONE CONSISTENTLY UNRELIABLE LINK, Raft can get into situations where LEADERSHIP CONTINUALLY BOUNCES BETWEEN TWO NODES, or the current leader is CONTINUALLY FORCED TO RESIGN, so the system EFFECTIVELY NEVER MAKES PROGRESS. (Addressed by a PRE-VOTE PHASE.) Paxos also depends on leaders and can have similar issues; EGALITARIAN PAXOS (EPaxos) uses a LEADERLESS protocol more robust against poorly performing nodes or connections.

On this page