Distributed Consensus

⭐ Interview Importance: MEDIUM
⏱️ Revision Time: 5 min

Concept

In a distributed system, how do you get 5 different servers to agree on the exact value of a variable, or agree on who the Leader is, when network messages can be delayed, dropped, or arrive out of order?
Distributed Consensus is the mathematical and programmatic solution to getting multiple independent nodes to agree on a single source of truth, even in the face of severe network failures.

The Two Generals’ Problem

To understand why consensus is so hard, consider the famous thought experiment:

  • Two Generals are on opposite hills, preparing to attack a city in the valley.
  • They must attack at the exact same time, or they will lose.
  • The only way to communicate is by sending a messenger running through the enemy valley.
  • General A sends a messenger: “Attack at dawn. Acknowledge this message.”
  • The messenger arrives. General B sends a reply: “Agreed. Attack at dawn.”
  • The Dilemma: General B now thinks, “What if my messenger was captured? If General A didn’t receive my reply, he won’t attack, and I will die. General A must send a confirmation that he received my reply.”
  • This creates an infinite loop of required confirmations. It is mathematically proven that perfect consensus over an unreliable network is impossible.

The Solution: Paxos and Raft

Since perfect consensus is impossible, computer scientists designed algorithms that achieve practical consensus using Majorities (Quorums) and Timeouts.

1. Paxos (1989)

The original consensus algorithm. It is notoriously difficult to understand and even harder to implement in code. It relies on a complex series of Proposers, Acceptors, and Learners passing multiple phases of messages to agree on a single value.

2. Raft (2013)

Created specifically because Paxos was too confusing. Raft achieves the exact same level of safety as Paxos but structures the logic much more clearly. It breaks consensus down into three sub-problems:

  1. Leader Election: (See previous topic).
  2. Log Replication: The Leader accepts writes from clients, appends them to its log, and sends them to the Followers.
  3. Safety: If a Leader dies, the algorithm guarantees that the newly elected Leader will be the node that has the most up-to-date log, preventing older nodes from overwriting new data.

Real-World Usage

You will almost never write a consensus algorithm yourself. You will use existing tools that implement them:

  • etcd: The backing store for Kubernetes (uses Raft).
  • Zookeeper: The classic coordinator for distributed systems (uses ZAB, a Paxos-like protocol).
  • Consul: HashiCorp’s service discovery tool (uses Raft).

Interview Questions

Q: Explain the difference between Crash Fault Tolerance (CFT) and Byzantine Fault Tolerance (BFT).
A:

  • Crash Fault Tolerance (Raft, Paxos): These algorithms assume that servers might crash or network cables might break, but they assume that if a server sends a message, that message is honest and correct. They are used in secure corporate networks (like AWS).
  • Byzantine Fault Tolerance (Blockchain): This assumes that nodes are actively malicious. A hacked node might send contradictory messages to different servers to intentionally crash the system. BFT algorithms (like Proof of Work or PBFT) are exponentially slower and more complex, required only in trustless environments like Cryptocurrency networks.

Q: Why do consensus systems (like Zookeeper) perform terribly if you deploy them across multiple global regions (e.g., US, Europe, Tokyo)?
A: Consensus algorithms require a majority of nodes to acknowledge every single write before the system can proceed. If you have nodes spread across the globe, every write requires network packets to cross the ocean, taking hundreds of milliseconds. This massive latency completely destroys the write-throughput of the system. Consensus clusters should generally be deployed within a single region (across multiple availability zones) to keep network latency under 5ms.