Leaderless Replication — N/R/W Quorums, Hinted Handoff, Read Repair & Anti-Entropy
Leaderless (Dynamo-style) replication uses N/R/W quorums, hinted handoff, read repair, and Merkle anti-entropy. R+W>N is not linearizability. Tombstones and clock-based LWW are the usual footguns.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Overview
Leaderless replication (from Amazon's Dynamo paper, used by Cassandra, ScyllaDB and Riak) has no leader to fail over. A coordinator sends each write to all N replicas of a key and reports success after W acks. A read asks R replicas and returns the newest version. If R + W > N, every read set overlaps every write set, so at least one replica in the read holds the latest acknowledged write. Replicas that missed writes are repaired three ways: read repair (fix stale replicas found during a read), hinted handoff (a stand-in node keeps writes for a down replica and replays them later) and anti-entropy (background Merkle-tree comparison). The catch is that quorums look strong but are not linearizable by default. Sloppy quorums, concurrent writes, LWW timestamps and partial write failures all leak anomalies.
Quorum math
Press Run. Snippets must be self-contained — no network, files, or native modules.
Output:
N R W R+W>N overlap-always stale-read-rate
3 2 2 True True 0.000
3 1 3 True True 0.000
3 3 1 True True 0.000
3 1 1 False False 0.670
5 2 2 False False 0.298
5 3 3 True True 0.000
5 1 2 False False 0.599| Setting (N=3) | Write tolerates | Read tolerates | Guarantee | Typical use |
|---|---|---|---|---|
W=2, R=2 (QUORUM/QUORUM) | 1 node down | 1 node down | Overlap: reads see the latest acked write | Default "strong-ish" |
| W=3, R=1 | 0 down for writes | 2 down for reads | Overlap; fast reads | Read-heavy, rare writes |
| W=1, R=3 | 2 down for writes | 0 down for reads | Overlap; fast writes | Write-heavy ingest |
W=1, R=1 (ONE/ONE) | 2 down | 2 down | No overlap: stale reads common | Metrics, logs, caches |
LOCAL_QUORUM (multi-DC) | Quorum in the local DC only | Same | Overlap within a DC; cross-DC async | Multi-region Cassandra |
Quorum write, quorum read, then repair
R+W>N is not linearizability. Repair and tombstones still matter.
- 1
Fan out the write
Send to all N replicas in the preference list. - 2
Wait for W acks
Sloppy quorum + hinted handoff covers down nodes. - 3
Read R replicas
Pick newest version; digests catch siblings. - 4
Repair in the background
Read repair and Merkle anti-entropy close the gaps.
Flow
- 1
1. Client write k equals v2
- next2. Coordinator finds N=3 replicas on the hash ring
- 2
2. Coordinator finds N=3 replicas on the hash ring
- nextReplica A: ack
- nextReplica B: ack
- downReplica C: unreachable
- 3. hint stored for CNeighbor D keeps the hint
- 3
Replica A: ack
- next4. W=2 acks so success
- 4
Replica B: ack
- next4. W=2 acks so success
- 5
Replica C: unreachable
- 6
Neighbor D keeps the hint
- 5. C returns: replay hintReplica C: unreachable
- 7
4. W=2 acks so success
- 8
6. Later read with R=2 sees v2 and v1
- next7. Return v2 and read-repair the stale replica
- 9
7. Return v2 and read-repair the stale replica
Lesson map
Leaderless Replication — N/R/W Quorums, Hinted Handoff, Read Repair & Anti-Entropy
Leaderless (Dynamo-style) replication uses N/R/W quorums, hinted handoff, read repair, and Merkle anti-entropy. R+W>N is not linearizability. Tombstones and clock-based LWW are the usual footguns.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB c["1. Client write k equals v2"] co["2. Coordinator finds N=3 replicas on the hash ring"] ok["4. W=2 acks so success"] rd["6. Later read with R=2 sees v2 and v1"] c -->|1. Client write k equals v2 to 2. Coordinator finds N=3 replicas on the hash ring| co
Repair mechanisms compared
| Mechanism | When it runs | Fixes | Cost | Gaps |
|---|---|---|---|---|
| Read repair | During a read that sees divergent versions | Keys that are actually read | Extra writes on the read path | Cold keys never get repaired |
| Hinted handoff | While a replica is down (bounded window, for example 3 h) | Writes missed during short outages | Hint storage on neighbors | Hints expire; the stand-in can fail too |
| Anti-entropy (Merkle trees) | Scheduled background repair (nodetool repair, Reaper) | Everything, eventually | CPU, IO and streaming | Must run more often than gc_grace_seconds, or deleted data comes back |
// Anti-entropy with Merkle trees: find which key ranges differ between two
// replicas by comparing hashes top-down instead of shipping every key.
// Sandbox-runnable: tsc --strict then node. No dependencies.
// 32-bit FNV-1a: fine for a demo; real systems use SHA-256 / xxHash / MD5.
const h = (s: string): string => {
let x = 0x811c9dc5;
for (let i = 0; i < s.length; i++) x = Math.imul(x ^ s.charCodeAt(i), 0x01000193) >>> 0;
return x.toString(16).padStart(8, "0");
};
type Store = Map<number, string>;
// Build leaf hashes over fixed key buckets, then hash pairs up to the root.
function tree(store: Store, buckets: number, keyspace: number): string[][] {
const per = keyspace / buckets;
const leaves: string[] = [];
for (let b = 0; b < buckets; b++) {
let acc = "";
for (let k = b * per; k < (b + 1) * per; k++) acc += `${k}=${store.get(k) ?? ""};`;
leaves.push(h(acc));
}
const levels = [leaves];
while (levels[levels.length - 1].length > 1) {
const prev = levels[levels.length - 1];
const next: string[] = [];
for (let i = 0; i < prev.length; i += 2) next.push(h(prev[i] + prev[i + 1]));
levels.push(next);
}
return levels; // levels[last][0] is the root
}
// Walk down from the root, descending only into subtrees whose hashes differ.
function diffBuckets(a: string[][], b: string[][]): { buckets: number[]; compared: number } {
let frontier = [0], compared = 0;
for (let lvl = a.length - 1; lvl >= 0; lvl--) {
const next: number[] = [];
for (const i of frontier) {
compared++;
if (a[lvl][i] !== b[lvl][i]) next.push(...(lvl === 0 ? [i] : [2 * i, 2 * i + 1]));
}
frontier = next;
}
return { buckets: frontier, compared };
}
const KEYS = 1024, BUCKETS = 16;
const r1: Store = new Map(), r2: Store = new Map();
for (let k = 0; k < KEYS; k++) { r1.set(k, `v${k}`); r2.set(k, `v${k}`); }
r2.set(70, "STALE"); r2.delete(900); // replica 2 missed two writes
const { buckets, compared } = diffBuckets(tree(r1, BUCKETS, KEYS), tree(r2, BUCKETS, KEYS));
console.log(`differing buckets: ${buckets.join(", ")} (keys ${buckets.map((b) => `${b * 64}-${b * 64 + 63}`).join(", ")})`);
console.log(`hash comparisons: ${compared} vs ${KEYS} key comparisons for a full scan`);Output:
differing buckets: 1, 14 (keys 64-127, 896-959)
hash comparisons: 15 vs 1024 key comparisons for a full scanWhy R + W > N is not linearizability
- Sloppy quorum: when home replicas are unreachable, writes go to other nodes (with hints). W acks no longer overlap the R home replicas, so the overlap guarantee is gone until the hints are replayed.
- Concurrent writes + LWW: two clients write different values, and the higher timestamp (subject to clock skew) wins on every replica. One acked write is silently discarded.
- Partial failure: a write reaches only 1 of the required 2 replicas and the client gets an error. The value is still on one replica and may later spread by read repair. A "failed" write becomes visible.
- Read concurrent with a write: one reader sees the new value and a later reader sees the old one, because the write hasn't reached a quorum yet. Fixing that needs blocking read repair and write-back before returning, and even then there are edge cases. For compare-and-set, Cassandra uses lightweight transactions (Paxos) at a much higher cost.
- Deletes: a delete is a tombstone. If a replica that missed the delete isn't repaired within
gc_grace_seconds, the tombstone is purged elsewhere and the old value comes back (zombie data).
What happens if you choose differently
- Leader-follower instead: you get a total order, transactions and simpler semantics, but a failover event and leader-bound writes. Choose it for relational, transactional data.
- ONE/ONE "for speed": reads are fast and almost always fine, until a node restarts and users see stale or missing data for minutes. Acceptable for telemetry, not for user-facing state.
- ALL for "safety": any single node down fails requests. You've built lower availability than a single leader.
- No scheduled repair: cold data drifts, and deleted data resurrects after the grace period. Running repair is not optional in Cassandra.
Pros and cons
Pros: no leader and no failover pause; tunable consistency per request; scales horizontally on a consistent-hash ring; tolerates node loss gracefully; strong write availability. Cons: weaker real semantics than the formula suggests; repair is an ongoing operational burden; tombstones and LWW lead to data loss and resurrection; no multi-key transactions; data modeling must be query-first.
Interview Q&A
Q1. N=3, W=2, R=2. One node is down. Can you read and write? Is the read always fresh? Yes to both operations, since 2 of 3 is available. A read overlaps the last acked write's replicas, so it sees it, assuming a strict (not sloppy) quorum and no concurrent writes. Stale replicas get read-repaired.
Q2. Explain hinted handoff and its limit. When a target replica is down, a neighbor stores the write as a hint and replays it when the replica comes back. That covers short outages. Hints have a time window and storage limit, and they don't replace anti-entropy repair. If the hinting node dies, the hint is lost.
Q3. Why do Cassandra operators fear gc_grace_seconds?
Tombstones are dropped after the grace period. A replica that missed the delete, and wasn't repaired within the period, still has the old value. Repair then copies that value back to everyone, and the deleted data resurrects. So the rule is: run a full repair more often than gc_grace_seconds (default 10 days).
Q4. How do Merkle trees make anti-entropy efficient? Each replica hashes its key ranges into a tree. Comparing roots tells you whether anything differs, and you descend only into differing subtrees, so network cost scales with the amount of difference, not the data size. Only the differing ranges are streamed.
Q5. When would you choose Cassandra/ScyllaDB over Postgres with replicas? For very high write throughput, data larger than one node, multi-DC active writes, simple access patterns known in advance (query-first tables), and an acceptance of eventual consistency or LWW. Choose Postgres for relational queries, transactions, constraints and ad-hoc analytics.
Go Deeper
- Amazon Dynamo paper (SOSP 2007)
- Apache Cassandra docs: Dynamo-style architecture
- Apache Cassandra docs: Hints
- Apache Cassandra docs: Read repair
- Jepsen: Cassandra analysis
Related
- Prev: Multi-Leader & Active-Active - Conflict Detection, LWW, CRDTs & Home Regions (
multi-leader-replication-conflicts-active-active) - Next: Database Replication for Engineers - Leader-Follower, Multi-Leader & Leaderless Quorums (
database-replication-leader-follower-multi-leader-leaderless)
This series:
- Database Replication for Engineers - Leader-Follower, Multi-Leader & Leaderless Quorums (
database-replication-leader-follower-multi-leader-leaderless) - Sync vs Async vs Semi-Sync Replication - Durability, RPO & Physical vs Logical Logs (
replication-sync-async-semi-sync-physical-logical) - Replication Lag & Session Guarantees - Read-Your-Writes, Monotonic Reads & Consistent Prefix (
replication-lag-read-your-writes-session-guarantees) - Failover & Split Brain - Detection, Promotion, Fencing & Lost Writes (
replication-failover-split-brain-fencing) - Multi-Leader & Active-Active - Conflict Detection, LWW, CRDTs & Home Regions (
multi-leader-replication-conflicts-active-active) - Leaderless Replication - N/R/W Quorums, Hinted Handoff, Read Repair & Anti-Entropy (
leaderless-replication-quorums-read-repair-anti-entropy) - this page
Existing Study pages (cross-link only, not rewritten here):
- Consistent hashing - rings, vnodes & replicas (
consistent-hashing-rings-vnodes-replicas) - CRDTs - conflict-free replicated data types (
crdts-conflict-free-replicated-data) - Database sharding - partition keys & rebalancing (
db-sharding-partitioning-keys-rebalancing)