yield point
Distributed Systemsintermediateupdated 2026-08-22

Quorums and R+W>N

One inequality decides whether your replicated store can return stale data, and it has nothing to do with how fast replication is.

Replicate your data across five machines. A write goes to some of them, a read comes from some of them. How many of each do you need before a read is guaranteed to see the most recent write?

The answer is a counting argument, and it is completely independent of network speed, replication lag, or clock accuracy.[1]paperWeighted Voting for Replicated DataGifford, D. K., SOSP '79, 1979

Set R=2, W=2 with N=5 and watch stale reads appear. Then kill replicas until writes stop entirely - that exact point is your availability limit.

R0v0R1v0R2v0R3v0R4v0
R + W vs N
3 + 3 = 6 > 5
Guaranteed overlap
1 replica(s)
Latest version
v0
Last read saw
v0
Stale reads
-
tick 0 / 1200
Break it
Do

Drop R+W to N or below and stale reads become possible.

Safety properties
  • ✓
    R+W>N means a read can never miss the latest write

    held on every tick so far

The pigeonhole

Pick any W replicas to write to and any R replicas to read from, out of N total. If R + W > N, those two sets cannot be disjoint - there are not enough replicas for them to avoid each other. At least R + W − N replicas appear in both.

Any replica in the intersection received the write, so the read sees it. The read takes the newest version among the replicas it consulted, and the newest version is guaranteed to be in there somewhere.

Notice what this argument does not rely on. Not timing. Not replication finishing. Not clocks. Just the sizes of two sets drawn from the same pool.

Set R=2, W=2 with N=5 in the widget. Now R+W = 4 ≤ 5, the sets can miss each other, and stale reads start appearing within a few rounds.

Choosing R and W

The inequality leaves you a family of choices, and they are genuinely different systems.

  • W=N, R=1 - writes must reach everyone, reads ask one replica. Fast reads, and any single failure blocks all writes.
  • R=N, W=1 - the mirror image. Writes never block; reads must reach every replica.
  • R=W=⌈(N+1)/2⌉ - the usual majority quorum. Balanced, and tolerates ⌊(N−1)/2⌋ failures on both paths.

Dynamo made these knobs explicit and per-operation, which is where the “tunable consistency” framing in Cassandra and Riak comes from.[2]paperDynamo: Amazon's Highly Available Key-value StoreDeCandia, G. et al., SOSP '07, 2007

Availability is what you pay

Every replica you require is a replica whose failure can block you. With N=5 and W=3, two failures are survivable and three are not.

Kill replicas one at a time in the widget and watch the write counter. There is a precise point where it stops, and that point is N − W + 1 failures.

This is also why a minority partition cannot make progress: it does not contain enough replicas to form a quorum, no matter how healthy those replicas are. Partition off three of five and the remaining two can neither read nor write at R=W=3.

What quorums do not give you

This is the part that catches people out. R+W>N guarantees a read sees the latest successfully completed write. It says nothing about:

Concurrent writes. Two clients writing different values to overlapping quorums at the same time both succeed. The intersection now holds two versions with no ordering between them, and something else - vector clocks, last-write-wins, application merge - has to decide. Dynamo returns both and makes the caller resolve it.[2]paperDynamo: Amazon's Highly Available Key-value StoreDeCandia, G. et al., SOSP '07, 2007

Failed writes. A write that reached some replicas but fewer than W is reported as failed and is not rolled back. A later read may see it. “Failed” means “not acknowledged”, not “did not happen”.

Atomicity across keys. Quorums are per key. There is no relationship between them.

Jepsen’s testing of Cassandra found exactly these gaps in practice - the quorum arithmetic was sound while the surrounding operations still lost updates.[3]articleJepsen: CassandraKingsbury, K., 2013 The inequality is necessary, not sufficient.

The dial

You gainRead-your-writes consistency without a leader, with tunable latency on each side
You payEvery replica you require is a replica whose failure blocks you

The line - when someone asks

In a quorum system a write must reach W of N replicas and a read must consult R of them, and if R+W>N the two sets are guaranteed to share at least one replica - which will hold the newest version. That is a pigeonhole argument, not a timing property, so it holds regardless of how slow replication is. The cost is availability: raising R and W tightens consistency and simultaneously increases the number of failures that can block you.

Recall

Loading…

Where are you with this?

Saved on this device. Sign in to keep it across devices.

Sources

Primary

Secondary