Consensus Algorithms: How Distributed Systems Agree on Truth

Consensus Algorithms: How Distributed Systems Agree on Truth 

A configuration management team at an infrastructure company once ran a three-node etcd cluster backing their service discovery system, and during a routine network maintenance window, a switch misconfiguration split the cluster into two isolated groups of one and two nodes.

The lone node, still receiving traffic from a subset of clients, kept accepting writes for several seconds before its client connections timed out. The two-node group, holding a majority, also kept operating.

When the network partition healed, the team had to reconcile two sets of writes that had happened during the split, a problem that etcd’s underlying consensus algorithm, Raft, was specifically designed to prevent in the first place, and would have, had that node not been reachable at all by any client during the partition.

Knowing exactly why the majority side stayed authoritative, and the minority side didn’t, is the core of what consensus algorithms solve. 

Two Nodes, One Leader, and a Network Partition 

Distributed systems need multiple machines to agree on a single, consistent version of the truth, which node is the current leader, what the latest committed value of a piece of data is, or whether a particular transaction happened, even when individual machines can fail, restart, or become temporarily unreachable from each other.

This sounds like it should be simple, but proving that a group of independent machines, communicating over an unreliable network, can reliably agree on anything is a deeply hard problem, one that took the distributed systems research community decades to solve rigorously. 

Consensus algorithms provide the formal machinery for this agreement, guaranteeing that as long as a sufficient number of nodes remain reachable and functioning correctly, the group will converge on one agreed value and never disagree about it, even in the presence of network delays, message loss, or node crashes.

Systems like etcd, Consul, ZooKeeper, and the internal coordination layers of most distributed databases all depend on a consensus algorithm underneath, whether or not the engineers using them think about it day to day. 

Quorums and the Majority Rule 

The foundational idea behind most practical consensus algorithms is the quorum: a decision is only considered final once a majority of nodes, more than half, have agreed to it. This single rule is what prevents the split-brain scenario from the introduction, because in any network partition, at most one side of the split can contain a majority of the total nodes, and only that side is allowed to keep making progress. 

For a cluster of five nodes, a quorum is three; the cluster can tolerate up to two node failures (or a partition isolating up to two nodes) and continue operating correctly, since the remaining three still form a majority.

This is why consensus clusters are almost always deployed with an odd number of nodes, three, five, seven, since an even-numbered cluster gains no additional fault tolerance over the next-smaller odd number while adding unnecessary coordination overhead, and worse, an even split (two and two, out of four) leaves neither side with a majority, halting the entire system rather than degrading gracefully. 

  • Quorum: the minimum number of nodes, typically a strict majority, required to agree before a decision is considered committed. 
  • Fault tolerance: for a cluster of 2f + 1 nodes, the system tolerates up to f node failures while remaining available. 
  • Odd cluster sizes: chosen deliberately to avoid ties during a network partition, since a tie leaves no side with a majority. 
  • Split-brain prevention: only the partition holding a quorum is allowed to accept new writes, keeping the system single-authority even during a split. 

It’s worth being precise about what a quorum really protects against, since it’s easy to overstate. Quorum-based consensus guarantees that at most one side of a partition can make progress, preventing two conflicting histories from both being accepted as valid. It does not, on its own, guarantee that the majority side is making the “right” decision in any business sense, only that whatever it decides becomes the single, agreed-upon truth the rest of the system can build on.

That distinction matters when reasoning about failure scenarios: a quorum cluster can be perfectly consistent internally while still reflecting a stale or incomplete view of the world if, for instance, the nodes that happen to form the majority were themselves cut off from some external source of information the minority side had access to. 

Raft: Leader Election and Log Replication 

Raft: Leader Election and Log Replication

Raft was designed explicitly to be easier to reason about than its predecessor, Paxos, without sacrificing the same correctness guarantees, and that focus on clarity is a large part of why it became the algorithm of choice for most new distributed systems built over the last decade, including etcd, Consul, and CockroachDB’s internal replication layer. 

Raft breaks consensus into two largely separate concerns: leader election and log replication. At any given time, the cluster has at most one leader, elected through a randomized-timeout mechanism where nodes that stop hearing from a leader after a timeout become candidates and request votes from their peers, with whichever candidate secures a majority of votes becoming the new leader. 

Node states: Follower -> Candidate -> Leader 

Follower: waits for heartbeat from leader; if timeout, 
           becomes Candidate and requests votes 
Candidate: requests votes from peers; becomes Leader if
           it wins a majority, or reverts to Follower 
Leader: accepts client writes, replicates them to
           followers, sends periodic heartbeats

Once elected, the leader is the only node that accepts client writes, appending them to its own replicated log and streaming that log to the follower nodes. A write is considered committed only once a majority of nodes have acknowledged appending it to their own log, at which point it’s guaranteed to survive even if the leader immediately crashes, since a majority of nodes already have a durable copy. 

  • Term: a logical clock incremented on every new election, used to detect and discard messages from an outdated, stale leader. 
  • Heartbeat: periodic messages a leader sends to followers, both to maintain its authority and to replicate new log entries. 
  • Log matching: Raft’s mechanism for ensuring followers’ logs stay consistent with the leader’s, resolving any divergence by having followers adopt the leader’s version. 
  • Randomized election timeout: staggering how long each follower waits before becoming a candidate, reducing the chance of repeated split votes among simultaneous candidates. 

The log replication half of Raft deserves its own attention, since it’s what really delivers durability once a leader is in place. Every client write is appended as a new entry to the leader’s log before anything else happens, tagged with the term in which it was created. The leader then sends that entry to every follower in parallel, and only once a majority (including the leader itself) has durably stored it does the leader consider the entry committed and respond to the client.

Followers that fall behind, because they were briefly disconnected, or restarted, catch up by receiving a batch of missed entries the next time they reconnect, with Raft’s log matching property guaranteeing that once two logs agree on an entry at a given index, they agree on every entry before it too, which makes catching up a matter of finding where the logs diverge and overwriting from that point forward rather than reasoning about the entire history from scratch. 

Paxos and Its Reputation for Complexity 

Paxos predates Raft by roughly two decades and remains foundational to the theory of distributed consensus, but it earned a reputation, even among experienced distributed systems engineers, for being notoriously difficult to fully grasp and correctly implement, a reputation Leslie Lamport himself acknowledged and which directly motivated Raft’s design. 

The core Paxos protocol (sometimes called “single-decree Paxos”) reaches agreement on a single value through a two-phase process, a prepare phase where a proposer gathers promises from a majority of acceptors, followed by an accept phase where the proposer’s value is formally agreed upon, provided no conflicting proposal has intervened.

Real systems need to agree on a continuous sequence of values, not just one, which requires “Multi-Paxos,” an extension that’s far less precisely specified in the original literature and has historically led to many subtly different, sometimes subtly incorrect, real-world implementations. 

Google’s Chubby lock service and Spanner’s internal replication both built on Paxos-family algorithms, proving the approach works at serious scale, but the operational and implementation complexity relative to Raft is a major reason most newer open-source distributed systems chose Raft instead when consensus became a widely needed building block rather than a specialized, research-heavy undertaking. 

Byzantine Fault Tolerance Considerations 

Raft and Paxos both assume what’s called the “crash-fault” model: a node either works correctly or stops entirely (crashes, becomes unreachable), but it never behaves maliciously or sends deliberately incorrect information while still appearing to participate normally.

This assumption holds well for most enterprise infrastructure, where nodes are controlled by the same organization and outright malicious behavior from your own hardware is not the primary threat model. 

Byzantine fault tolerance addresses a stronger, more adversarial threat model, where nodes might actively lie, send conflicting information to different peers, or otherwise behave arbitrarily, whether due to a bug, corruption, or outright malicious intent.

This is essential in contexts like public blockchains, where nodes are run by mutually distrusting parties with no shared organizational accountability, but it’s a heavy, unnecessary cost for a typical internal distributed system where crash-fault tolerance is already sufficient protection. 

  • Crash-fault tolerance: assumes failed nodes simply stop, the model Raft and Paxos both operate under. 
  • Byzantine fault tolerance: assumes failed nodes can behave arbitrarily or maliciously, requiring stronger, more expensive protocols like PBFT. 
  • Higher node count requirements: Byzantine fault tolerance typically requires 3f + 1 nodes to tolerate f faulty ones, versus 2f + 1 for crash-fault tolerance. 
  • Use case separation: internal infrastructure almost always uses crash-fault-tolerant algorithms; public, trust-minimized systems generally require Byzantine fault tolerance. 

Consensus in Practice: etcd, ZooKeeper, and Consul

These three tools are the most common places engineers directly encounter consensus algorithms in production, typically without needing to implement the algorithm themselves, since each tool exposes a simpler interface (a key-value store, a lock service, a service registry) built on top of its internal consensus layer. 

etcd, built on Raft and most famous as Kubernetes’s backing store for cluster state, exposes a straightforward key-value API while handling leader election, log replication, and failure recovery internally.

ZooKeeper, older and built on its own consensus protocol called Zab (conceptually similar to Raft in its leader-based approach), remains widely deployed in Kafka’s older architecture (prior to Kafka’s own move toward KRaft, its own internal Raft-based mode) and in various Hadoop ecosystem tools.

Consul, also Raft-based, focuses more specifically on service discovery and health checking, layering that functionality on top of the same underlying consensus guarantees. 

  • etcd: Raft-based, key-value oriented, central to Kubernetes control plane state.
  • ZooKeeper: Zab-based, widely used historically for coordination in the Hadoop and older Kafka ecosystems. 
  • Consul: Raft-based, specialized for service discovery, health checking, and configuration.
  • KRaft: Kafka’s newer, self-contained Raft-based consensus layer, removing its historical dependency on ZooKeeper entirely. 

Latency and Availability Trade-Offs 

Latency and Availability Trade-Offs

Consensus isn’t free, and the cost shows up directly in write latency, since every write has to be replicated to and acknowledged by a majority of nodes before it’s considered durable, rather than committing on a single machine the way an unreplicated database might. This makes consensus-backed systems a poor fit for extremely high-throughput, latency-sensitive write workloads unless that durability guarantee is worth the cost, which for coordination and configuration data, exactly the use case etcd, ZooKeeper, and Consul target, it usually is. 

Geographic distribution amplifies this cost further: a five-node cluster spread across three continents for disaster-tolerance purposes pays a substantial latency penalty on every write, since the leader has to wait for round-trip acknowledgment from enough geographically distant followers to form a majority.

Many teams instead cluster consensus nodes within a single region or a small number of nearby regions, accepting a smaller blast radius for a regional outage in exchange for consensus latency that stays in the single-digit milliseconds rather than the hundreds of milliseconds a global spread would introduce on every single write. 

A common compromise pattern separates the two concerns: run the consensus cluster itself within one region for low-latency writes, and layer a separate, asynchronous replication mechanism to other regions on top for disaster recovery, accepting that the disaster-recovery copy trails behind by some bounded amount rather than being part of the strongly consistent quorum.

This mirrors the broader pattern from database replication more generally, synchronous coordination where consistency truly matters, asynchronous replication for the reach that would otherwise be too costly to make synchronous at all.

Teams evaluating a consensus-backed system for the first time often assume the tool should own their full global footprint end to end, when in practice it usually makes more sense as the strongly consistent core of a system whose outer layers use looser guarantees deliberately. 

Final Thoughts 

Consensus algorithms answer a question that sounds almost philosophical, how do independent machines agree on the truth when any of them might fail at any moment, with mathematically rigorous, provably correct machinery that real infrastructure depends on constantly, usually invisibly.

Raft’s leader election and log replication, built on the simple but powerful idea of requiring majority agreement, is why a Kubernetes cluster’s state stays coherent even as nodes come and go, and why the infrastructure team’s split etcd cluster in the opening story had a well-defined, correct answer about which side was allowed to keep operating.

The algorithm itself is rarely something application engineers need to implement from scratch, but knowing roughly how it works, why odd cluster sizes matter, why a minority partition stops accepting writes, why consensus adds latency proportional to geographic spread, turns what would otherwise be a mysterious outage into a predictable, well-explained consequence of a design decision made deliberately, for good reason.

Frequently Asked Questions 

Can a consensus cluster survive losing more than half its nodes? 

Not while continuing to accept new writes safely, that’s the fundamental trade-off consensus algorithms make. Losing a majority means no side of any resulting partition has enough nodes to form a quorum, so the system halts writes to avoid the risk of split-brain, prioritizing consistency over availability during that window. 

Is Raft used only for small clusters? 

Raft clusters are typically kept small (three, five, sometimes seven nodes) because every additional node adds replication overhead without adding proportional fault tolerance benefit. Systems needing to scale beyond that typically shard data across many independent, smaller Raft groups rather than growing a single group very large. 

How does consensus differ from eventual consistency? 

Consensus algorithms provide strong consistency for the specific decisions they govern, everyone agrees on the same value, immediately and permanently, once it’s committed. Eventual consistency, common in systems like Cassandra, allows temporary disagreement between nodes that resolves over time, trading immediate agreement for higher availability and lower latency. 

Do all distributed databases use a consensus algorithm internally? 

No. Many distributed databases use leaderless replication with quorum reads and writes (a related but distinct approach) rather than a full consensus algorithm like Raft or Paxos, especially systems designed to prioritize write availability during partitions over strict consistency. 

What happens during a leader election in Raft if two candidates request votes simultaneously? 

A split vote can occur where no candidate secures a majority, in which case the election times out and a new one begins after each candidate’s randomized timeout expires, with the randomization specifically intended to reduce the chance of another perfectly simultaneous, deadlocked split vote. 

Is Byzantine fault tolerance ever used outside of blockchains? 

Occasionally, in specific high-assurance contexts like aerospace or certain financial systems where the cost of a compromised or malfunctioning node behaving deceptively is severe enough to justify the added overhead, but it remains rare in mainstream backend infrastructure compared to crash-fault-tolerant algorithms.

Similar Posts