Distributed Consensus

Consensus is how a group of machines agrees on a value, or a sequence of values, even when some machines crash and messages are delayed. It's the foundation of reliable distributed systems: electing a single leader, keeping replicas of a database log identical, storing cluster configuration, and implementing locks that really are exclusive. Without consensus, network hiccups lead to split brain: two nodes both believing they're primary and accepting conflicting writes.

In practice, you rarely implement consensus yourself. You rely on systems built on Raft or Paxos: etcd (behind Kubernetes), ZooKeeper, Consul, CockroachDB, TiKV, Kafka's KRaft mode, and Spanner. Understanding how they work explains their guarantees, their failure behavior, and why they need an odd number of nodes.

TL;DR

Quick Example

Using etcd (Raft-based) for leader election among service instances:

Only one instance holds leadership at a time. If it crashes or loses its lease, another takes over within seconds.

Core Concepts

The Consensus Problem

A consensus protocol must satisfy:

For replicated services, the usual goal is state machine replication: all replicas apply the same commands in the same order, so they stay identical.

FLP Impossibility

The Fischer-Lynch-Paterson result (1985) proves that in a fully asynchronous network, where messages can be delayed arbitrarily, no deterministic protocol can guarantee consensus if even one node may crash. Practical protocols keep safety (never disagreeing) always, and guarantee liveness (making progress) only when the network behaves reasonably, using timeouts and randomization.

Raft

Raft splits consensus into understandable parts:

  1. Leader election: nodes are followers, candidates, or the leader. If followers hear no heartbeat within a randomized timeout, one becomes a candidate, increments the term, and requests votes. A candidate with votes from a majority becomes leader. Randomized timeouts avoid repeated split votes.
  2. Log replication: clients send commands to the leader, which appends them to its log and replicates them to followers. Once a majority has stored an entry, it's committed and applied to the state machine.
  3. Safety: a node only votes for candidates whose log is at least as up to date as its own, so a new leader always has every committed entry. Leaders never overwrite committed entries.

Extras include log compaction via snapshots, joint consensus or single-server changes for membership changes, and read optimizations (ReadIndex, leader leases) for linearizable reads without writing to the log.

Paxos

Lamport's Paxos (1989/1998) is the original practical consensus algorithm: proposers, acceptors, and learners run prepare and accept phases with majority quorums. Multi-Paxos elects a stable leader to skip the first phase for sequences of values. It's powerful, but notoriously hard to understand and implement fully, which motivated Raft. Google's Chubby and Spanner use Paxos variants.

Quorums and Fault Tolerance

Majority quorums guarantee that any two quorums overlap, so a new leader always learns committed decisions.

Use odd sizes. Spread nodes across failure domains (three availability zones for three nodes), and keep latency between them low, since every commit waits for a majority round trip.

Leases and Fencing Tokens

Leader leases and distributed locks built on consensus have a subtle hazard: a leader can pause (a GC pause or VM freeze) after its lease expires, then wake up believing it's still leader and write to storage. Fencing tokens fix this: each new leader or lock holder gets a monotonically increasing number (an etcd revision, a ZooKeeper zxid), and storage systems reject writes carrying a token older than the highest seen.

Byzantine Fault Tolerance

Raft and Paxos assume nodes fail by crashing, not by lying. Byzantine fault tolerant protocols (PBFT, Tendermint, HotStuff) tolerate malicious nodes, at the cost of 3f + 1 nodes and more messages. They're used in blockchains and some high-assurance systems. See blockchain fundamentals.

Systems Built on Consensus

Best Practices

Use a Proven Implementation

Consensus is easy to get subtly wrong. Use etcd, ZooKeeper, Consul, or battle-tested libraries (etcd/raft, hashicorp/raft), rather than writing your own for production.

Keep Consensus Clusters Small and Close

Three or five nodes is typical. More nodes increase write latency and message overhead. Place them in separate zones of one region; cross-region consensus adds tens to hundreds of milliseconds per commit.

Store Small, Critical Data Only

Coordination stores like etcd and ZooKeeper are for metadata, configuration, leader election, and locks, not for bulk application data. Watch their size limits and write rates.

Always Fence

Any leader-based or lock-based design touching external resources should pass fencing tokens and have those resources enforce them. Leases alone aren't enough.

Common Mistakes

Even-Sized or Two-Node Clusters

Two nodes tolerate zero failures (a majority of 2 is 2), and four tolerate the same one failure as three. Use 3 or 5.

Assuming Locks Are Exclusive Without Fencing

Fencing tokens checked by the storage layer prevent the stale write.

Treating Consensus Stores as General Databases

Putting large blobs or high-volume writes into etcd degrades the whole cluster, including Kubernetes if it's the cluster's etcd. Use it for small, low-churn coordination data.

FAQ

Why do consensus clusters need an odd number of nodes?

Progress requires a majority. With 2f + 1 nodes you tolerate f failures, and adding one more node (an even count) raises the majority size without tolerating any additional failures. It only adds cost and latency.

What's the difference between Raft and Paxos?

Both solve consensus with majority quorums and provide the same core guarantees. Raft was designed for understandability: a strong leader, explicit terms, and clear rules for elections, log replication, and membership changes. Paxos is more general and minimal, but harder to implement completely. Most new systems choose Raft.

What happens during a network partition?

The side with a majority can elect a leader and keep committing writes. The minority side can't reach a quorum, so it stops accepting writes (and linearizable reads) until the partition heals. That's consensus choosing consistency over availability.

Does Kubernetes use consensus?

Yes. Kubernetes stores all cluster state in etcd, which uses Raft. The API server reads and writes through etcd, so etcd's availability, which needs a majority of its members, is critical to the control plane.

Related Topics

References