Quorums and leaderless replication
How Dynamo-style databases replicate without a leader: N, W and R quorums, what quorum overlap guarantees and what it does not, sloppy quorums and hinted handoff, read repair, anti-entropy with Merkle trees, conflict handling with vector clocks or last-writer-wins.
Reading is half of it. See this used in a real interview: walk through Design a Distributed Key-Value Store →
Amazon's Dynamo paper introduced a style of replication without a leader: any replica can accept writes, and clients read and write to several replicas at once. Cassandra, ScyllaDB, Riak and DynamoDB's ancestors use it. "Design a key-value store" usually expects this design, with its vocabulary of quorums, hinted handoff, read repair and Merkle trees. This guide explains each and what guarantees they really give.
The setup
- Data is partitioned with consistent hashing; each key is stored on N replicas (the next N nodes on the ring).
- A write is sent to all N replicas and succeeds when W acknowledge.
- A read queries replicas and returns when R respond, picking the newest version among them.
No leader means no failover pause: as long as enough replicas respond, the system keeps working.
Quorums
If W + R > N, every read set overlaps every write set in at least one replica, so a read should see the latest successful write. Common choices with N = 3:
| W | R | Behaviour |
|---|---|---|
| 2 | 2 | balanced; tolerates one node down for both reads and writes |
| 3 | 1 | fast reads, but writes fail if any replica is down |
| 1 | 3 | fast writes, slow reads |
| 1 | 1 | fastest, eventually consistent, may read stale data |
Lower W and R mean lower latency and higher availability; higher values mean stronger consistency. Many systems let you choose per request. See CAP theorem.
What quorums do not guarantee
Even with W + R > N, edge cases can return stale data or lose ordering:
- Concurrent writes to the same key may land in different orders on different replicas.
- A write that succeeded on fewer than W replicas (reported as failed) may still be visible on some.
- Sloppy quorums (below) break the overlap guarantee.
- Clock-based conflict resolution can discard newer writes.
Quorums give "probably the latest", not linearizability. For strict guarantees (compare-and-set, uniqueness), use consensus-based systems. See consensus and coordination.
Sloppy quorums and hinted handoff
If some of a key's N home replicas are down, a strict quorum would fail the write. A sloppy quorum writes to other healthy nodes instead, each storing a hint that the data belongs elsewhere. When the home replica recovers, the hint is handed off to it. Writes stay available during failures, at the cost of temporarily weaker consistency.
Repairing divergence
Replicas drift apart through failures and missed writes. Three mechanisms bring them back:
- Read repair: when a read sees replicas with different versions, it writes the newest version back to the stale ones.
- Hinted handoff: replays missed writes after a short outage.
- Anti-entropy: a background process compares replicas and copies missing data. Comparing every key would be too expensive, so replicas build Merkle trees (hash trees over key ranges): comparing root hashes finds whether ranges differ, and walking down finds exactly which keys differ, exchanging little data.
Resolving conflicts
When concurrent writes produce different values:
- Last-writer-wins (LWW) with timestamps: simple, but concurrent writes are silently lost and clock skew can let older writes win. Cassandra uses this by default.
- Vector clocks or version vectors: detect that two versions are concurrent rather than ordered, keeping both as siblings for the application (or a merge function) to resolve. Dynamo's shopping cart merged siblings by union.
- CRDTs: data types that merge automatically (counters, sets). See offline-first apps and sync.
See unique ids, ordering and time for clocks.
Failure detection and membership
Nodes learn who is alive and which ranges they own through gossip, and decide when to route around a node with failure detectors. See gossip and failure detection.
Leader-based versus leaderless
| Leader-based | Leaderless | |
|---|---|---|
| Writes go to | the leader | any W replicas |
| Failover | elect a new leader, short pause | none needed |
| Ordering | leader orders writes | concurrent writes need resolution |
| Consistency | strong reads from the leader possible | tunable, eventual by default |
| Latency under failures | spikes during failover | stable, tail latency from slowest of W or R |
See replication and consistency.
In the interview
For Design a Key-Value Store: consistent hashing with virtual nodes, N = 3 replicas, tunable W and R (2 and 2 by default), sloppy quorums with hinted handoff, read repair, Merkle-tree anti-entropy, gossip membership, and LWW or vector clocks for conflicts, with an honest statement of what consistency that gives.
Checklist
- N replicas per key on the ring; W and R chosen per use case.
- W + R > N for read-your-latest in normal operation, with known edge cases.
- Sloppy quorums and hinted handoff for availability.
- Read repair and Merkle-tree anti-entropy for convergence.
- Conflict policy: LWW, vector clocks with siblings, or CRDTs.
- Consensus-based storage where strict guarantees are required.