Consistency modelsQuorum reads and writes

100%

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.

Intermediate19 minUpdated 30 Sept 2026

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

Replicas of cart:91 (N = 3)
  1. A, 1 item, In the write set only, value v7
  2. B, 1 item, In both sets (the overlap), value v7
  3. 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.
Start

As it starts. 4 steps follow.

Three replicas of cart:91. The first step is R = 2, W = 2: the write reached A and B, the read asked B and C, and B, in both sets, carries v7 into the read. Step through the other settings; the last one is R = 1, W = 1, where the sets can miss each other and the read returns the old v6.

N = 3: failures each path survives

N, R, WReads survive (N − R down)Writes survive (N − W down)R + W > N?With one replica down
3, 1, 1reads: 2 downwrites: 2 downNo — 2 ≤ 3Both paths keep working; a read may return v6 after v7 was acknowledged
3, 1, 3reads: 2 downwrites: 0 downYes — 4 > 3Every write fails until the replica returns
3, 2, 2reads: 1 downwrites: 1 downYes — 4 > 3Both paths keep working; a second failure stops both (the Dynamo paper's common setting)
3, 3, 1reads: 0 downwrites: 2 downYes — 4 > 3Every read fails until the replica returns
Was this section helpful?

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

A W = 2 write, then an R = 2 read, as an ordered list of steps:
A W = 2 write, then an R = 2 read17 steps between Client, Coordinator, Replica A, Replica B, Replica C. The steps are listed as text after the diagram.Replica CReplica BReplica ACoordinatorClientW = 2 reached; stop waitingslow disk: still at v6highest version wins: v7put(cart:91, v7), W = 21store v72store v73store v74ack5ack6ok7get(cart:91), R = 28read9read10v611v712read repair v713v714
  1. Client → Coordinator: put(cart:91, v7), W = 2
  2. Coordinator → Replica A: store v7
  3. Coordinator → Replica B: store v7
  4. Coordinator → Replica C: store v7
  5. Replica A → Coordinator (reply): ack
  6. Replica B → Coordinator (reply): ack
  7. Note over Coordinator: W = 2 reached; stop waiting
  8. Coordinator → Client (reply): ok
  9. Note over Replica C: slow disk: still at v6
  10. Client → Coordinator: get(cart:91), R = 2
  11. Coordinator → Replica B: read
  12. Coordinator → Replica C: read
  13. Replica C → Coordinator (reply): v6
  14. Replica B → Coordinator (reply): v7
  15. Note over Coordinator: highest version wins: v7
  16. Coordinator → Replica C: read repair v7
  17. 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, WR + W > N?Read survivesWrite survivesFeels like
5, 3, 3Yes — 6 > 5reads: 2 downwrites: 2 downMajority both ways; any two writes also meet
5, 1, 5Yes — 6 > 5reads: 4 downwrites: 0 downRead-heavy data that rarely changes
5, 2, 4Yes — 6 > 5reads: 3 downwrites: 1 downLeans toward reads, still tolerates a write-side failure
5, 2, 2No — 4 ≤ 5reads: 3 downwrites: 3 downPartial 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.

Sloppy quorums
When a home replica is down, Dynamo sends the write to the next healthy node on the ring and counts its ack towards W. That stand-in keeps a hint and hands the value home later (hinted handoff). Until it does, a read quorum of the home replicas can miss the write even though R + W > N. Cassandra takes the stricter line: a hint never counts towards the consistency level, except at ANY.
Partial writes
A write that reached 1 of 3 replicas and then timed out is reported to the client as failed, but it was not rolled back. A later read that happens to include that replica returns it, and read repair spreads it, so a 'failed' write may or may not survive. Retries must be idempotent for exactly this reason.
Concurrent writes and clocks
Highest version wins only if versions order writes correctly. With last-writer-wins on client timestamps, a client whose clock runs 2 s fast beats a genuinely later write, and the overlap replica faithfully returns the wrong winner. Vector clocks keep both versions instead; the versioning topic covers them.
Reads that go backwards
While a write is still in flight, one reader can see the new value on the replica that has it, and a later reader can ask two replicas that do not and see the old one. Linearizable quorum reads prevent this by writing the value back to a quorum before returning (the ABD algorithm). Cassandra does this by default: with read_repair = BLOCKING, a QUORUM read that finds disagreeing replicas waits until its repair writes reach the consistency level, which the docs call monotonic quorum reads. For compare-and-set it adds Paxos-based lightweight transactions.

Two reads during one slow write

Scenario 1 of 2: As described.

Notes
  • 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.

  1. 0–47 ms · Writer · put(v7), W = 2 (ok) — Returns at 47 ms, once the second ack (A's) arrives.
  2. 1 ms · Writer → Replica B, arriving 10 ms
  3. 1 ms · Writer → Replica A, arriving 45 ms: slow link
  4. 1 ms · Writer → Replica C, arriving 50 ms
  5. 10 ms · Replica B · has v7
  6. 10–45 ms · Replica A, Replica B, Replica C · window: only B has v7
  7. 12–26 ms · Reader 1 · get → v7 (B, C) (ok)
  8. 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.
  9. 45 ms · Replica A · has v7
  10. 50 ms · Replica C · has v7
N = 3, R = 2, W = 2, times illustrative. The write of v7 reaches B at 10 ms but A only at 45 ms and C at 50 ms. Reader 1 asks B and C and gets v7; reader 2 starts after reader 1 has finished, asks A and C, and gets v6. The second tab adds ABD's write-back: reader 1 writes v7 to C before returning, so reader 2 cannot miss it.
Was this section helpful?

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

Assumptions
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)
Working
  1. 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
  2. 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
  3. 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
  4. 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
What it means
  • 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

Assumptions
Each replica is up
99% of the timeillustrative; failures assumed independent
Minutes in a year
525,600
Working
  1. N = 3, at least 1 up (R or W = 1)1 − 0.01³0.999999from Each replica is up
  2. N = 3, at least 2 up (majority)3 × 0.99² × 0.01 + 0.99³0.999702from Each replica is up
  3. N = 3, all 3 up (W = 3)0.99³0.970299from Each replica is up
  4. 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
  5. N = 5, all 5 up0.99⁵0.950990from Each replica is up
  6. 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
  7. 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
What it means
  • 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

Minutes a year with no quorumRequiring every replica is the worst choice: all 5 of 5 is down about 25,800 minutes a year, a majority of 5 about 5.0.11101001k10k100k5 of 53 of 34 of 52 of 33 of 51 of 325.8k min15.6k min515 min157 min5.2 min0.53 minQuorumMinutes per year unavailable (min)Minutes a year with no quorumRequiring every replica is the worst choice: all 5 of 5 is down about 25,800 minutes a year, a majority of 5 about 5.0.11101001k10k100k5 of 53 of 34 of 52 of 33 of 51 of 325.8k min15.6k min515 min157 min5.2 min0.53 minQuorumMinutes per year unavailable (min)
Each replica up 99% of the time, independently (an illustrative assumption). A write that needs W acks fails when fewer than W replicas are up; the same numbers apply to a read with R. Log scale.
Data
Quorumminutes per year
5 of 525,760
3 of 315,611
4 of 5515
2 of 3157
3 of 55.2
1 of 30.53

What real stores do

Amazon Dynamo
The paper gives (N, R, W) = (3, 2, 2) as the common setting. Writes use a sloppy quorum over the first N healthy nodes, with hinted handoff to return data home, so the cart stays writable during failures. Services that must never reject a write can set W = 1.
Apache Cassandra
QUORUM is ⌊RF/2⌋ + 1 replicas, so 2 of 3 or 3 of 5. With two datacentres at RF 3 each, QUORUM counts all 6 and needs 4; LOCAL_QUORUM needs 2 of the 3 in the coordinator's datacentre; EACH_QUORUM needs 2 of 3 in both. By default (read_repair = BLOCKING) a QUORUM read that finds replicas disagreeing waits for its repair writes to reach the level before answering, so quorum reads never go backwards.
Azure Cosmos DB
Each region keeps a set of four replicas. At every level a write commits to a local majority of 3; strong and bounded-staleness reads read a local minority of 2. 3 + 2 > 4, so within a region this is an exact quorum pair running in production; weaker levels read one replica. For strong consistency on a multi-region account, the write must also commit in every region (a global majority).
Measured staleness (PBS)
Bailis and colleagues modelled partial quorums (R + W ≤ N) with production latency traces and found reads usually consistent within tens of milliseconds of a write. That is why many teams run R = W = 1 on purpose for data where a brief stale read is harmless.
Was this section helpful?

Trade-offs.

The chosen option is first; the others stay visible so the reasoning can be checked.

01
A cart service on N = 3
Chosen:R = 2, W = 2
  • 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
Downside we accept:
  • Con:Every request waits for the second-fastest reply (9 ms in the example, not 4)
  • Con:Two replicas down stops both paths
Ruled out:R = 1, W = 1

Reads can miss acknowledged writes (a partial quorum); Concurrent writes can land on different replicas and conflict silently

Ruled out:R = 1, W = 3

Any single replica down or slow blocks every write (97% write availability at 99% per replica)

02
Across two datacentres
Chosen:LOCAL_QUORUM for reads and writes
  • Pro:No request waits for the cross-datacentre round trip
  • Pro:Each datacentre keeps serving on its own if the link between them fails
Downside we accept:
  • 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
Ruled out:EACH_QUORUM writes

Every write pays the cross-datacentre round trip; A failed link or datacentre blocks all writes

Ruled out:QUORUM across both datacentres

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

FailureImpactDetectionMitigationMeanwhile
A replica returns after a long outage5Replica CIt serves values that are hours old to R = 1 reads that land on itRead-repair rate and hint backlog jump for that nodeReplay hints, and run a full Merkle-tree repair before it serves R = 1 traffic; R ≥ 2 reads repair it as they goReads with R + W > N still return the newest version
The coordinator crashes after one ack2CoordinatorThe client sees a timeout, yet the write exists on one replica and may surface laterClient timeouts without a matching error from the storeRetry with the same version or idempotency key, so a duplicate is harmlessThe 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 itMonitor stored hints per node and their ageUse a strict quorum where freshness matters; replay or expire hints before they pile upWrites stay available; reads may be stale
Clock skew under last-writer-winsA newer write carries an older timestamp and silently losesRarely detected; shows up as lost updates reported by usersVersions assigned by the coordinator or the replica, vector clocks, or conditional writes (lightweight transactions); the strong-vs-eventual topic covers last-writer-wins versus mergingNo error is ever raised
Was this section helpful?
Builds on this
Replication and quorums
Read next