Consider a key replicated on 5 nodes (N=5) with write quorum W=3 and read quorum R=2. Walk through why a read is, or is not, guaranteed to return the most recently acknowledged write. How would the answer change if R were increased to 3?
answer
- boundary case R+W=N fails
- 3+2=5 no overlap forced
- raise R to 3 -> 3+3=6>5
- pigeonhole again
- latency/availability cost of raising R
basics
~20 sWith 5 copies, writing to 3 and reading from only 2 isn't enough — 3+2 just equals 5, so the 2 you read from could be exactly the 2 that were skipped by the write. Bump the read to 3 copies and now 3+3=6 is more than 5, so the two groups must share at least one copy.
solid answer
~40 sR+W must exceed N for the overlap guarantee to hold; here R+W=3+2=5=N, so it's on the boundary and fails. The 5 replicas can be split into the 3 that got the write and the other 2 — and if the read happens to query exactly those other 2, it will observe none of the replicas that hold the latest version, returning stale data with no way to know it's stale. Raising R to 3 makes R+W=6>5=N: any 3-of-5 read set and any 3-of-5 write set must intersect, because two disjoint subsets of a 5-element set can have at most 5 elements combined, and 3+3=6 exceeds that. So R=3 restores the guarantee at the cost of a slower read (waiting on one more replica).
go deeper
Should recognize that 3+2=5 is not greater than N=5 and therefore something is off, even if they can't fully derive the disjoint-set argument.
Should walk through the concrete disjoint-partition counterexample and correctly compute that R=3 fixes it, tying it back to R+W>N.
Should immediately flag the boundary condition as a common misconfiguration, and discuss the latency/availability cost of the fix, not just that it works.
Should connect this to real operational incidents (e.g. ONE/ONE consistency misconfigurations) and discuss alternative fixes (raise W instead, add session consistency, use read repair) with their respective trade-offs.
## Why the boundary case fails This scenario is a deliberately boundary case, because R+W=N (rather than R+W>N) is a common off-by-one mistake when people reason about quorums quickly, and understanding exactly why it fails clarifies the whole mechanism. Start with the write: - The coordinator sends the new value to some or all of the 5 replicas and waits until 3 of them acknowledge durable storage before telling the client the write succeeded. - Which 3 replicas actually respond first is not fixed in advance — it depends on network latency and load at that moment, so over many writes, different 3-of-5 subsets can end up holding the 'latest' write while the remaining 2 are still on an older version (they'll eventually catch up via replication, but at the moment right after the write, they may not have it yet). ## The read that can miss Now the read: the coordinator queries **2 of the 5 replicas** and returns whichever response looks newest. The critical question is whether those 2 replicas are guaranteed to include at least one of the 3 that held the latest write. With 5 total replicas split into a hypothetical 3 (wrote) and 2 (didn't write yet), the read's 2-replica query could, in the worst case, land exactly on those 2 laggards — a perfectly valid subset of size 2 out of 5, disjoint from the write's 3. In that case, both replicas the read consults are still on the old version, and the client returns stale data with full confidence, because from its point of view both replicas it asked agreed (there's no third opinion contradicting them). This isn't a rare pathological case — under normal operation, if the read coordinator happens to prefer contacting the fastest-responding or geographically nearest replicas, and those overlap poorly with which replicas happened to ack the write fastest, stale reads become a real, recurring behavior, not just a theoretical edge case. ## The counting argument The reason the arithmetic matters is the **pigeonhole principle**: two subsets of an N-element set can only be guaranteed to intersect if their sizes sum to more than N. 1. At exactly R+W=N, you can always construct a partition (write-set, read-set) that exactly tiles the N replicas with zero overlap — that's precisely what happened above (3 + 2 = 5, tiling all 5 replicas exactly). 2. Only once R+W > N does the counting argument break: if the write's 3 replicas and the read's (now) 3 replicas were disjoint, that would require 3+3=6 distinct replicas, but only 5 exist, which is impossible. 3. So at least one replica must be counted in both sets — guaranteeing the read touches at least one replica holding the latest write. ## What raising the read quorum costs Raising R from 2 to 3 has real operational cost, which is the trade-off every quorum tuning decision makes. - **Latency.** The read now has to wait for 3 replicas to respond instead of 2, so its latency is bounded by the third-fastest replica rather than the second-fastest — typically a modest but real increase, and it becomes worse under replica slowness or partial outages. - **Fault tolerance.** Since now 3 of the 5 replicas must be reachable for the read to even complete (versus 2 before), reducing the operation's fault tolerance from 3 simultaneous failures survivable down to 2. This is the general pattern: pushing the consistency guarantee (via R+W>N) costs latency and availability headroom on whichever side you increase. ## Where this shows up in production A concrete production analog: this is exactly the class of bug that shows up when a team configures a Cassandra-style store with, say, ONE-consistency writes and ONE-consistency reads against RF=3 (N=3, effectively W=1,R=1, R+W=2<3) expecting 'good enough' consistency, then is surprised by visibly stale reads right after a write under load — and the fix discussion is identical to this scenario: - raise R - raise W - accept the eventual-consistency window and add client-side techniques like read-your-writes session stickiness to the same replica that served the write
- Instead of raising R, could you fix this by raising W to 4 while keeping R=2?Yes — 4+2=6>5 also restores the overlap guarantee. The choice between raising R or raising W depends on whether you'd rather slow down writes or reads; raising W makes every write wait for one more replica instead.
- Does the overlap guarantee tell you which of the R replicas returned the freshest data, or just that one of them did?Just that at least one did — the client still needs a way to identify which response is newest, typically a version number, vector clock, or timestamp attached to each replica's value, and the client compares the R responses and picks the highest one.
- If the read coordinator always queries the two geographically closest replicas rather than a random two, does that change the guarantee?The mathematical guarantee is unaffected — it only depends on set sizes, not which specific replicas are chosen. But practically, a biased selection strategy (always the same 'closest' replicas) can make stale reads either more or less likely to manifest in practice depending on which replicas tend to lag, even when R+W ultimately does satisfy the quorum condition.
Picture 5 mailboxes. A courier drops the memo into 3 of them (the write). If you only check 2 mailboxes at random, you might unluckily check exactly the 2 that weren't touched. Check a 3rd mailbox and now it's mathematically impossible for your 3 checked boxes and the courier's 3 dropped boxes to miss each other entirely, since there are only 5 boxes total.
saying these in an interview costs you the question
- Says R+W=N is fine because it's 'close enough'
- Can't explain why the specific two skipped replicas are the failure case
- Thinks raising R or W has no cost
- Confuses this with leader election or consensus quorums (majority voting for leadership) rather than data freshness overlap