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 Data
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 Store
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 Store
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: Cassandra The inequality is necessary, not sufficient.
The dial
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
- [1]
- [2]
Secondary
- [3]