SysDesignPrep.com
Study guide 66 of 183

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:

WRBehaviour
22balanced; tolerates one node down for both reads and writes
31fast reads, but writes fail if any replica is down
13fast writes, slow reads
11fastest, 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-basedLeaderless
Writes go tothe leaderany W replicas
Failoverelect a new leader, short pausenone needed
Orderingleader orders writesconcurrent writes need resolution
Consistencystrong reads from the leader possibletunable, eventual by default
Latency under failuresspikes during failoverstable, 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.

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.