Consensus, Clocks and Coordination
Know how nodes agree despite crashes, partitions and unreliable clocks: majority quorums, leader terms, leases with fencing, logical clocks and when a distributed transaction blocks.
Key points
- 1
Consensus (Raft, Paxos) needs a majority of floor(N/2) + 1. A cluster of 2f + 1 members tolerates f failures, so an even member count adds cost without adding tolerance.
- 2
Leaders are identified by a term or epoch. A deposed leader can still be alive, so linearizable reads must confirm leadership with a quorum (ReadIndex) or a time-bounded lease.
- 3
Never trust a lock alone for correctness. Pair it with a monotonically increasing fencing token that the protected resource checks and enforces.
- 4
Measure durations with a monotonic clock. Order events across machines with Lamport or vector clocks, or with bounded-uncertainty time such as TrueTime (commit wait ≈ 2ε).
- 5
2PC is atomic but blocks if the coordinator fails after participants vote yes. Sagas avoid distributed locks but give up isolation and need idempotent, compensating steps.
Common traps
A lower Lamport timestamp does not prove one event caused the other; only vector clocks can detect concurrency.
Checking that you still hold a lock right before writing does not help: a pause can happen between the check and the write.
Spreading 5 nodes as 2 / 2 / 1 over three zones survives a zone loss, but not a zone loss plus one more failure.
Read the source
Test yourself on Consensus, Clocks and Coordination
Ten questions, with the answer and explanation after each one.