HotShard
Replication & consistency

Quorum Replication

Tune the consistency and availability trade-off with three numbers: N, R, and W.

Copy a value to several servers and two questions appear at once. When is a write done? And which servers does a read have to ask? Quorum replication answers both with three numbers: N copies, W acks per write, R replies per read. Dynamo, Cassandra, and Riak run on it. One rule, R + W > N, decides whether a read is guaranteed to overlap the latest write. This page shows the rule, what it buys, and where it stops short.

~8 min read

Start here: the problem it solves#

TL;DRthe 30-second version
  • Leaderless replication stores every value on N nodes. A write waits for W acks. A read collects R replies. No node is special, so there is no leader to fail.
  • If R + W > N, the read set and the write set share at least one replica. So a read always sees a copy of the latest successful write.
  • Among the R replies, the newest version wins. Lagging replicas are patched by read-repair on the read path and by anti-entropy in the background.
  • R + W > N is necessary but not sufficient for linearizability. Clock skew, sloppy quorums, and concurrent writes can still lose or hide data.

Picture a value you want to survive machine failures, so you copy it to several servers. Now two questions have no obvious answer. When is a write done? And which servers does a read have to ask? If every replica must confirm every write, you are unavailable the moment one node is slow or rebooting. If a read asks just one replica, it is fast, but it can hand back data that is seconds or even hours out of date.

Single-leader systems answer this by sending every write through one elected node. Raft does this, and so does a MySQL or Postgres primary. That gives a clean order. But it brings back a special node. When it dies, writes stall until a new leader is elected. Leaderless replication throws the leader away. Every replica accepts reads and writes directly. Clients talk to several replicas at once and reconcile what they hear.

Quorum replication turns 'how many replicas must agree' into a dial. You pick how many acks a write needs (W) and how many replies a read needs (R), out of N copies. This is what tunable consistency means. The same database can serve a strongly consistent read, and then an eventually consistent one on the next request, just by changing R and W.

The core trade-offBigger R and W give stronger consistency, but more replicas must be reachable and must answer before the operation returns. Smaller R and W give fast, highly available operations that may read stale data. Because there is no leader, the system keeps serving as long as a quorum is reachable.

The mechanism: N replicas, W acks, R reads#

Three numbers define the system. N is the replication factor: how many copies of each key exist. In Cassandra those are the next N nodes around a consistent-hashing ring. W is the write quorum: how many replicas must ack before the coordinator tells the client 'committed'. R is the read quorum: how many replicas the coordinator must hear from before it answers. The coordinator sends every request to all N replicas and waits for the first W or R to respond.

  1. A write arrives. The coordinator sends the new value, tagged with a version, to all N replicas in parallel.
  2. Each reachable replica stores the value and replies with an ack. As soon as W acks arrive, the coordinator returns success. It does not wait for the slow or down replicas.
  3. A read arrives. The coordinator asks all N replicas in parallel and waits for the fastest R to reply.
  4. The coordinator compares the R replies by version and returns the newest. If the versions disagree, it has found divergence.
  5. Read-repair: the coordinator pushes the newest version back to any replica that returned a stale or missing value, nudging the system back toward agreement.

How does the read know which reply is newest? Every write carries a version. Cassandra uses a wall-clock microsecond timestamp. Riak and classic Dynamo use a version vector. The read takes the highest version. With timestamps that is a simple last-write-wins comparison. With version vectors the read can tell that two writes were concurrent, and it hands both back as siblings for the application to merge.

W=2 (any 2 nodes must ack a write) Β· R=2 (any 2 must answer a read) Β· highlighted = the overlap

Av2write ackin the write quorum
Bv2write ack + readin BOTH quorums β€” the overlap
Cv1lagging, but readin the read quorum
N = 3, W = 2, R = 2 β€” why the read always overlaps the write

Why must they overlap? The write landed on at least W replicas. The read heard from at least R. Both sets come from the same N replicas. If they shared no node, together they would hold at least R + W distinct replicas. But R + W is bigger than N. So they must share at least one node, and that node holds the fresh value. Two things to notice. This only holds for a read that starts after the write returned success. And the read still has to pick the fresh value out by version. That second step is where the guarantee gets shaky, and the trade-offs section covers it.

PredictN = 5, W = 3. A write is acknowledged by 3 replicas, then two of those three crash before anything else happens. A reader uses R = 3. Is it guaranteed to see the new value?

Hint: How many surviving replicas still hold the new value, and how many distinct replicas can R = 3 reach?

Yes, but notice why, because the real danger is elsewhere. Only one surviving replica still holds the new value, and only three replicas are alive at all. A read of R = 3 must contact all three survivors, including that one. So newest-wins returns the fresh value. The catch is durability, not staleness. The acked write now lives on a single node. If it crashes too, the write is gone despite being 'committed'. Quorum systems acknowledge before the value is fully replicated, so a fresh write can still be lost to a burst of correlated failures.

Latency, failure tolerance, and common configurations#

The cost of a quorum operation is set by tail latency, not average latency. A write does not return until the W-th fastest replica responds. A read waits for the R-th fastest. So raising W or R samples further into the slow tail. One straggler, such as a garbage-collection pause or a busy disk, can dominate. This is why W = N or R = N is fragile: one slow replica stalls every operation. Smaller quorums ignore the slowest N βˆ’ W or N βˆ’ R replicas.

  • Failure tolerance: a write succeeds as long as W replicas are reachable, so it tolerates N βˆ’ W failures. Reads tolerate N βˆ’ R. With a majority quorum, both tolerate ⌊(Nβˆ’1)/2βŒ‹ failures.
  • Latency: write latency is the W-th fastest replica's round-trip; read latency is the R-th fastest. Lower R and W mean lower latency and more straggler tolerance, at the cost of freshness.
  • Storage: every value is stored N times no matter what R and W are. They only govern how many copies you wait for.
ConfigurationR + W vs NGood forCost
W = N, R = 1= N + 1 > NRead-heavy: reads hit any single replica, very fast and always freshWrites need every replica up β€” fragile, high write latency
W = 1, R = N= N + 1 > NWrite-heavy: writes ack from one replica, very fastReads need every replica up β€” fragile, high read latency
W = R = ⌈(N+1)/2βŒ‰> N (by 1)Balanced: both are a majority; tolerates the most node failures symmetricallyBoth reads and writes pay majority latency
W = R = 1 (N = 3)= 2 ≀ 3Maximum availability / lowest latencyStale reads allowed β€” eventual consistency only

The majority quorum is the usual default: 2 of 3, or 3 of 5. R + W beats N by exactly one, so both quorums stay as small as they can. That is Cassandra's QUORUM level and Dynamo's typical (3, 2, 2).

Trade-offs: why R + W > N is necessary but not sufficient#

It is tempting to read the overlap rule as 'R + W > N gives strong consistency, done.' That overclaims. The rule guarantees that a quorum read overlaps the latest committed write. That is enough for read-your-writes freshness in the simple, no-failures case. It does not guarantee linearizability: the property that the whole system behaves as if every operation happened at one instant, in one global order that every client agrees on.

  • Last-write-wins by wall-clock: if conflicts are settled by physical timestamps, clock skew can make an older write win over a newer one. The fresh value is among the replies, but the read picks the wrong one. Acknowledged writes get silently dropped.
  • Read-repair races: one read can return the new value while a concurrent read returns the old one, because read-repair has not reached the lagging replica yet. Two readers disagree about the order of events.
  • Sloppy quorums break overlap outright. During a partition, Dynamo lets a write collect its W acks from any reachable nodes, not just the key's home replicas. The stand-in stores the value with a hint and hands it to the home node when it recovers (hinted handoff). Write availability goes up. But a later strict read of the home nodes need not intersect that write.
  • Concurrent writes need conflict resolution. Two clients writing the same key through different coordinators can both reach W replicas without seeing each other. Without version vectors or CRDTs this becomes a lost update or a resurrected delete.
CAP and PACELC framingUnder a partition you choose: refuse operations that can't reach a quorum (lean CP, lose availability) or accept them through a sloppy quorum (lean AP, lose consistency). PACELC adds the part CAP omits: with no partition, you still trade latency for consistency. Every step up in R or W buys consistency by spending latency. No setting escapes the trade; you only move along the curve.
Comparison: quorum vs single-leader vs multi-leader
Quorum (leaderless)Single-leader (Raft)Multi-leader
Who accepts writesAny replica; coordinator fans out to NOnly the elected leaderSeveral leaders, one per region
Write availabilityHigh β€” any W reachable replicas, no electionStalls during leader election after a leader failureHigh β€” each region's leader is independent
ConsistencyTunable via R, W; linearizable only with careLinearizable by construction (single ordered log)Eventual; conflicts between regions need resolution
Conflict handlingVersions / vector clocks / LWW; siblings or read-repairNone needed β€” leader serializes all writesRequired β€” concurrent regional writes conflict
Typical systemsDynamo, Cassandra, Riak, Voldemort, DynamoDBetcd, ZooKeeper, CockroachDB ranges, Postgres primaryActive-active MySQL, BDR, CouchDB, multi-region DynamoDB GT

A single leader gives you a clean global order for free, but one point of write coordination and an election gap when it dies. Leaderless quorums have no election gap and no write bottleneck, but must rebuild enough ordering after the fact, through versions and overlap. Multi-leader sits between: fast regional writes, always paying for cross-region conflict resolution.

Where quorum replication runs in the wild
  • Amazon Dynamo (2007): the SOSP paper that defined the design. Consistent-hashing placement, (N, R, W) quorums, sloppy quorums with hinted handoff, vector clocks, and Merkle-tree anti-entropy. Everything below copies it.
  • Apache Cassandra: per-query consistency levels (ONE, QUORUM, LOCAL_QUORUM, EACH_QUORUM, ALL). Writes always go to all replicas; the level only sets how many acks the coordinator waits for. LOCAL_QUORUM needs a majority only in the client's datacenter, so a read doesn't wait on a replica across an ocean. Conflicts are last-write-wins by timestamp.
  • Riak: closest to classic Dynamo. Per-bucket N, R, W, version vectors that surface siblings, and CRDT data types (counters, sets, maps) so merges are automatic.
  • Amazon DynamoDB: hides N, R, W behind two read modes, eventually consistent (cheaper, may be stale) and strongly consistent. Writes use a quorum across availability zones.
  • Project Voldemort (LinkedIn): an open-source Dynamo clone with the same N, R, W knobs.

These systems agree on the quorum mechanics and differ on how they resolve concurrent writes. Cassandra picks last-write-wins (simple, can lose updates under clock skew). Riak keeps siblings or uses CRDTs (safe merges, more application work). DynamoDB papers over it with its two-mode read API.

Common mistakes
  • Treating R + W > N as linearizability. Jepsen showed Cassandra losing about 28% of acknowledged writes under QUORUM, even with synchronized clocks, because last-write-wins by timestamp kept the wrong version.
  • Calling a stale read under R = W = 1 a bug. With R + W ≀ N the two sets can be disjoint, so a read can hit only replicas that missed the write. That is the trade-off you chose, not a defect.
  • Trusting a W-acked write to be durable. The W replicas that acked may not have propagated to the other N βˆ’ W yet. If they fail first, the 'committed' write is gone. Higher W shrinks this window but never closes it.
  • Deleted values coming back. A delete is a tombstone: a versioned 'this key is gone' marker. If read-repair or anti-entropy copies an old value onto a replica before the tombstone reaches it, and the tombstone is later garbage-collected, the deleted value is resurrected.
  • Picking an even N. A majority of 4 tolerates 1 failure, the same as N = 3, for one extra replica of cost. Odd N gives the best failure tolerance per replica.
In an interview

Lead with the dial: N replicas, W acks per write, R replies per read, and the overlap rule R + W > N for read-your-writes freshness. Draw the N=3, W=2, R=2 picture and say the pigeonhole argument out loud: any two writers and any two readers must share a node. Name the systems: Dynamo, Cassandra (ONE / QUORUM / LOCAL_QUORUM / ALL), Riak, DynamoDB.

Then show senior judgment by stating the limit: R + W > N is necessary but not sufficient for linearizability. Mention clock skew under last-write-wins, sloppy quorums breaking overlap, and concurrent writes needing version vectors or CRDTs. Citing Jepsen's lost writes in Cassandra under QUORUM shows you know the gap between the formula and reality.

Be ready to reason about specific settings: why W=N, R=1 gives fast fresh reads but fragile writes, why the majority quorum tolerates the most failures, and why one straggler hurts a large quorum. Name the healing mechanisms: read-repair on the read path, anti-entropy with Merkle trees in the background, hinted handoff during failures.

Does R + W > N give me linearizability?

No. It only guarantees that a quorum read overlaps the latest committed write. Linearizability also needs correct ordering of concurrent writes (version vectors or consensus, not skewed wall-clocks), strict rather than sloppy quorums, and care around read-repair races. Necessary, not sufficient.

What's a sloppy quorum?

A quorum that takes its W acks from any reachable nodes, not just the key's home replicas. A write that can't reach W home nodes lands on stand-ins, which store it with a hint and forward it once the home nodes recover. It boosts availability during partitions. The cost is that a later strict read need not overlap the write.

What breaks on concurrent writes?

Two writes to the same key can each reach W replicas without seeing each other. Last-write-wins keeps only one, a lost update. Version vectors detect the concurrency and keep both as siblings for the app to merge. CRDTs make the merge automatic and order-independent.

Why is an odd N recommended?

A majority quorum on odd N tolerates ⌊(Nβˆ’1)/2βŒ‹ failures while keeping R + W just one above N. Even N buys the same tolerance as the next-lower odd number for one extra replica, and can create ties.

References & further reading
References

Feedback on this topic β†’