Quorum reads and writes.
With N copies of every key you don't have to wait for all of them. Wait for W on a write and R on a read, and if the two groups must overlap, the read is guaranteed to touch a replica with the newest acknowledged value. Moving R and W trades read speed against write speed, and fault tolerance against freshness, one request at a time.
Builds on Strong and eventual consistency.
The idea.
Three copies of a shopping cart, and a rule for how many of them each read and each write must hear from.
Why the overlap is guaranteed
A key-value store keeps the cart cart:91 on N = 3 replicas, A, B and C. Waiting for all three on every request would make each request as slow as the slowest machine, and one dead machine would stop everything. So a write waits for W = 2 acknowledgements and a read asks R = 2 replicas.
Count the slots. The write set and the read set together name 2 + 2 = 4 replicas, but there are only 3 to choose from, so at least one replica is in both. That replica holds the new version, v7. Every value carries a version, and the read returns the highest one it sees, so it returns v7 even if its other reply is an old v6. This is the pigeonhole principle, and it holds whichever two replicas each side happens to pick.
The rest of the topic is choosing N, R and W, and knowing exactly what the rule does not promise.
Which replicas each request touches
- A, 1 item, In the write set only, value v7
- B, 1 item, In both sets (the overlap), value v7
- C, 1 item, In the read set only, value v6, Answers the read with v6; the coordinator ignores it because B's v7 is newer.
- In the write set only
- In the read set only
- In both sets (the overlap)
- Read it and got the old value
- Not asked
- C
- Answers the read with v6; the coordinator ignores it because B's v7 is newer.
As it starts. 4 steps follow.
N = 3: failures each path survives
| N, R, W | Reads survive (N − R down) | Writes survive (N − W down) | R + W > N? | With one replica down |
|---|---|---|---|---|
| 3, 1, 1 | reads: 2 down | writes: 2 down | No — 2 ≤ 3 | Both paths keep working; a read may return v6 after v7 was acknowledged |
| 3, 1, 3 | reads: 2 down | writes: 0 down | Yes — 4 > 3 | Every write fails until the replica returns |
| 3, 2, 2 | reads: 1 down | writes: 1 down | Yes — 4 > 3 | Both paths keep working; a second failure stops both (the Dynamo paper's common setting) |
| 3, 3, 1 | reads: 0 down | writes: 2 down | Yes — 4 > 3 | Every read fails until the replica returns |
How it works.
One write and one read of cart:91 through a coordinator, then the settings for five replicas, then the four ways the overlap rule falls short.
A W = 2 write, then an R = 2 read
- Client → Coordinator: put(cart:91, v7), W = 2
- Coordinator → Replica A: store v7
- Coordinator → Replica B: store v7
- Coordinator → Replica C: store v7
- Replica A → Coordinator (reply): ack
- Replica B → Coordinator (reply): ack
- Note over Coordinator: W = 2 reached; stop waiting
- Coordinator → Client (reply): ok
- Note over Replica C: slow disk: still at v6
- Client → Coordinator: get(cart:91), R = 2
- Coordinator → Replica B: read
- Coordinator → Replica C: read
- Replica C → Coordinator (reply): v6
- Replica B → Coordinator (reply): v7
- Note over Coordinator: highest version wins: v7
- Coordinator → Replica C: read repair v7
- Coordinator → Client (reply): v7
W is how many acks to wait for, not how many copies to make
The coordinator sends the write to all N replicas every time; W only decides when it stops waiting and answers the client (Cassandra documents exactly this). The third copy still arrives, usually a few milliseconds later. Likewise a read may ask only R replicas, or ask all N and answer after the first R replies. When the replies disagree, the coordinator returns the newest and writes it back to the replica that was behind: read repair. Stragglers that nobody reads are fixed later by background anti-entropy that compares Merkle trees, which the key-value-store topic covers.
Settings for N = 5
| N, R, W | R + W > N? | Read survives | Write survives | Feels like |
|---|---|---|---|---|
| 5, 3, 3 | Yes — 6 > 5 | reads: 2 down | writes: 2 down | Majority both ways; any two writes also meet |
| 5, 1, 5 | Yes — 6 > 5 | reads: 4 down | writes: 0 down | Read-heavy data that rarely changes |
| 5, 2, 4 | Yes — 6 > 5 | reads: 3 down | writes: 1 down | Leans toward reads, still tolerates a write-side failure |
| 5, 2, 2 | No — 4 ≤ 5 | reads: 3 down | writes: 3 down | Partial quorum: fast, may be stale, writes can conflict |
Why W > N/2 matters on its own
R + W > N is about reads meeting writes. A second rule is about writes meeting each other: if W > N/2, any two write sets overlap, because two disjoint majorities cannot fit in N replicas. Two clients that update cart:91 at the same moment therefore collide on at least one replica, so a store that versions writes (vector clocks, or a conditional write) can see both and keep them as siblings or refuse the stale one. With disjoint write sets, no replica ever holds both. A plain last-writer-wins store gains nothing from the overlap: each replica just applies whichever write arrives. Vogels puts it the other way round: with W < (N + 1)/2 there is a possibility of conflicting writes that no replica ever sees side by side.
Majorities are also what Raft and Paxos use to commit a log entry, for the same reason: any two majorities share a member, so a new leader always finds every committed entry. This topic does not teach consensus; it is enough to recognise that a quorum store with W = R = majority and a consensus log rest on the same arithmetic.
What R + W > N does not promise
The overlap argument assumes a fixed set of N home replicas, a write that either finished or never happened, and versions that order writes correctly. Each assumption can break.
Two reads during one slow write
Scenario 1 of 2: As described.
- 1 Returns at 47 ms, once the second ack (A's) arrives.
- 8 Started after reader 1 returned v7, yet sees v6: the value went backwards. Each read met the rule; the pair is not linearizable.
Timeline as a list
Two reads during one slow write: 6 lanes, from 0 ms to 60 ms.
- 0–47 ms · Writer · put(v7), W = 2 (ok) — Returns at 47 ms, once the second ack (A's) arrives.
- 1 ms · Writer → Replica B, arriving 10 ms
- 1 ms · Writer → Replica A, arriving 45 ms: slow link
- 1 ms · Writer → Replica C, arriving 50 ms
- 10 ms · Replica B · has v7
- 10–45 ms · Replica A, Replica B, Replica C · window: only B has v7
- 12–26 ms · Reader 1 · get → v7 (B, C) (ok)
- 28–42 ms · Reader 2 · get → v6 (A, C) (violation) — Started after reader 1 returned v7, yet sees v6: the value went backwards. Each read met the rule; the pair is not linearizable.
- 45 ms · Replica A · has v7
- 50 ms · Replica C · has v7
In practice.
What a quorum costs in milliseconds and in downtime, and the settings real stores ship with.
Quorum latency is the k-th fastest reply
- Replica A replies after
- 4 mssame zone as the coordinator (illustrative)
- Replica B replies after
- 9 msanother zone (illustrative)
- Replica C replies after
- 60 msfar zone or mid garbage-collection pause (illustrative)
- W = 1 (or R = 1) waits for the fastest replymin(4, 9, 60)4 msfrom Replica A replies after, Replica B replies after and Replica C replies after
- W = 2 waits for the second-fastest reply2nd of (4, 9, 60)9 msfrom Replica A replies after, Replica B replies after and Replica C replies after
- W = 3 waits for the slowest replymax(4, 9, 60)60 msfrom Replica A replies after, Replica B replies after and Replica C replies after
- Cost of waiting for all three instead of two60 ÷ 9≈ 6.7× slowerfrom W = 2 waits for the second-fastest reply and W = 3 waits for the slowest reply
- Waiting for all N makes every request as slow as the slowest replica, and some replica is always slow. Dynamo notes that latency is dictated by the slowest of the R (or W) replicas, which is why R and W are kept below N.
- A majority hides one slow replica completely; that alone often justifies N = 3 over N = 2.
How often a quorum can be formed
- Each replica is up
- 99% of the timeillustrative; failures assumed independent
- Minutes in a year
- 525,600
- N = 3, at least 1 up (R or W = 1)1 − 0.01³0.999999from Each replica is up
- N = 3, at least 2 up (majority)3 × 0.99² × 0.01 + 0.99³0.999702from Each replica is up
- N = 3, all 3 up (W = 3)0.99³0.970299from Each replica is up
- N = 5, at least 3 up (majority)0.99⁵ + 5 × 0.99⁴ × 0.01 + 10 × 0.99³ × 0.01²0.999990from Each replica is up
- N = 5, all 5 up0.99⁵0.950990from Each replica is up
- W = 3 on N = 3: minutes a year writes are blocked(1 − 0.970299) × 525,600≈ 15,600 min (10.8 days)from N = 3, all 3 up (W = 3) and Minutes in a year
- Majority on N = 3: minutes a year with no quorum(1 − 0.999702) × 525,600≈ 157 minfrom N = 3, at least 2 up (majority) and Minutes in a year
- W = N turns three 99% machines into a 97% write path, worse than any one machine; a majority of three gives about 99.97%.
- Five replicas with majority quorums reach about 99.999% while still tolerating two failures.
- Replicas in one zone fail together (power, network, a bad deploy), so independence only holds if the replicas are spread across zones.
Minutes a year with no quorum
Data
| Quorum | minutes per year |
|---|---|
| 5 of 5 | 25,760 |
| 3 of 3 | 15,611 |
| 4 of 5 | 515 |
| 2 of 3 | 157 |
| 3 of 5 | 5.2 |
| 1 of 3 | 0.53 |
What real stores do
Trade-offs.
The chosen option is first; the others stay visible so the reasoning can be checked.
- Pro:Tolerates one replica down on both reads and writes
- Pro:A read sees every write acknowledged before it began
- Pro:Majority writes make concurrent updates meet on a replica
- Con:Every request waits for the second-fastest reply (9 ms in the example, not 4)
- Con:Two replicas down stops both paths
Reads can miss acknowledged writes (a partial quorum); Concurrent writes can land on different replicas and conflict silently
Any single replica down or slow blocks every write (97% write availability at 99% per replica)
- Pro:No request waits for the cross-datacentre round trip
- Pro:Each datacentre keeps serving on its own if the link between them fails
- Con:A read in the other datacentre may be stale until replication arrives
- Con:Losing a whole datacentre can lose writes it had not yet shipped
Every write pays the cross-datacentre round trip; A failed link or datacentre blocks all writes
With RF 3 + 3, a quorum is 4, so every request needs at least one remote reply; Losing either datacentre leaves 3 of 6 replicas, too few for any request
What goes wrong
| Failure | Impact | Detection | Mitigation | Meanwhile |
|---|---|---|---|---|
| A replica returns after a long outage5Replica C | It serves values that are hours old to R = 1 reads that land on it | Read-repair rate and hint backlog jump for that node | Replay hints, and run a full Merkle-tree repair before it serves R = 1 traffic; R ≥ 2 reads repair it as they go | Reads with R + W > N still return the newest version |
| The coordinator crashes after one ack2Coordinator | The client sees a timeout, yet the write exists on one replica and may surface later | Client timeouts without a matching error from the store | Retry with the same version or idempotency key, so a duplicate is harmless | The retry, not the lost reply, decides the outcome |
| Hinted-handoff backlog grows (sloppy quorum) | A write counted as a success is missing on its home replicas, so reads miss it | Monitor stored hints per node and their age | Use a strict quorum where freshness matters; replay or expire hints before they pile up | Writes stay available; reads may be stale |
| Clock skew under last-writer-wins | A newer write carries an older timestamp and silently loses | Rarely detected; shows up as lost updates reported by users | Versions assigned by the coordinator or the replica, vector clocks, or conditional writes (lightweight transactions); the strong-vs-eventual topic covers last-writer-wins versus merging | No error is ever raised |