Distributed systems
Quorum reads and writes: interview questions and how to answer them
A quorum system replicates to N nodes and requires W acknowledgements to write and R responses to read; when R + W > N the read set and the write set must overlap, so a read sees at least one copy of the latest completed write.
Written and reviewed by Sahil Srivastav
What it actually is
Quorum replication makes consistency a tunable rather than a fixed property. With N replicas, a write is acknowledged once W of them have stored it and a read waits for R responses. The arithmetic that matters is R + W > N: any two sets of that size must intersect, so a read always touches at least one replica holding the most recent completed write.
The overlap guarantee is precise and narrower than people assume. It means the newest completed write is *present in the response set*, so a reader that can identify the newest version returns it — which requires versioning, because otherwise the reader has several values and no way to order them. Quorum without version metadata does not give you the latest value; it gives you a set containing it.
It also says nothing about writes still in flight. A write that has reached two of three replicas but not yet been acknowledged may or may not be visible, and may later be seen and then not seen, because a failed write is not rolled back. That is why a strict quorum is not linearizability, only a strong-enough-for-most-purposes staleness bound.
Why it matters in production
Because it is the dial that turns a replicated store from fast-and-stale to slow-and-fresh without changing the data model. `W=1, R=1` gives low latency and unbounded staleness; `W=N, R=1` makes reads cheap and writes fragile to a single slow replica; `W=quorum, R=quorum` is the balanced default nearly every production Cassandra or Dynamo-style deployment runs.
The choice also decides the failure tolerance. With N=3 and quorum on both sides, the system survives one replica loss for both reads and writes. With W=3 it survives none for writes, which is how well-intentioned “make it consistent” configuration changes cause availability incidents during routine node replacement.
And it is where interviewers test whether a candidate understands that a formula is not a guarantee. Reciting R + W > N is easy. Explaining why it does not prevent reading a value that later disappears, and what you add to fix that, is the answer that lands.
How it works
Why overlap works
Pigeonhole. If a completed write is on W replicas and a read consults R replicas out of the same N, then R + W > N forces at least one replica to be in both sets. So the newest committed value is in the read response. The coordinator then resolves the responses — highest version, or latest timestamp, or a merge — and returns one answer.
Versioning is mandatory, not optional
The coordinator needs an order over the returned values. Cassandra uses per-cell timestamps and last-write-wins, which is simple and loses concurrent updates. Dynamo uses vector clocks and returns siblings when versions are concurrent, pushing the merge to the application. Either way, the reconciliation rule is a design decision with data-loss consequences, and it is separate from the quorum arithmetic.
Partial writes are not rolled back
If W=2 of N=3 is required and only one replica stores the value before the coordinator gives up, the client gets an error but the value is on one replica. A subsequent quorum read may or may not include it, and read repair may then propagate it, so a “failed” write can become permanent later. This is the monotonicity violation that makes quorum reads non-linearizable.
Sloppy quorums break the arithmetic on purpose
Under partition, Dynamo-style systems will accept W acknowledgements from *any* W reachable nodes, not the N home replicas, storing hints to be handed off later. This preserves availability but means the read set and write set may not intersect at all, so the R + W > N guarantee no longer holds during the partition. Knowing that the arithmetic assumes a strict quorum is the distinguishing detail.
Anti-entropy behind the quorum
Quorum reads bound staleness for keys that are read; keys that are never read stay divergent. Read repair fixes replicas lazily when a mismatch is observed, and Merkle-tree comparison in a repair process fixes the rest in the background. Without the background process, a replica that missed writes while down stays wrong for cold data indefinitely.
Implementing it
Default to N=3 with quorum reads and quorum writes, spread across failure domains. It tolerates one replica loss on both paths and gives a clear staleness story, which is a better starting point than tuning per query.
If a specific operation needs real linearizability — a uniqueness check, a balance guard — use the store’s consensus path (Cassandra lightweight transactions, DynamoDB conditional writes and transactions) rather than tightening R and W. Quorum arithmetic cannot provide compare-and-set.
Never set W=N to “be safe”. It makes every write depend on the slowest replica and on all of them being up, converting a tail-latency problem into an availability problem.
Keep repair running and monitor its progress. Teams often discover their anti-entropy job has been failing for weeks only when a replica replacement surfaces data that should have been deleted.
Interview questions and how to answer them
What does R + W > N guarantee?
That the read set and write set overlap, so a read consults at least one replica holding the most recent *completed* write. Given version metadata, the coordinator can then return that value. It does not guarantee linearizability: a write in progress may be visible then invisible, because a quorum write that fails is not rolled back.
N=3, W=2, R=2. A write succeeds on two replicas and the third is down. What can a reader see?
Any two of the three replicas, and at least one of them has the new value, so the reader sees it — that is the overlap property. When the third replica returns, it is stale until read repair or anti-entropy updates it, and a read that reaches it plus one up-to-date replica still returns the new version because the coordinator picks the higher version.
Why is a quorum read not linearizable?
Because concurrent and failed writes are not atomic across replicas. A write that reached one replica but failed to reach a quorum can be observed by a read that happens to include that replica, and then not observed by the next read that does not — a value appearing and disappearing, which linearizability forbids. Cassandra’s `LOCAL_SERIAL`/`SERIAL` paths add Paxos for exactly this reason.
What is hinted handoff and how does it interact with the quorum guarantee?
When a home replica is unreachable, the coordinator stores the write on another node as a hint and replays it when the replica returns. It keeps writes accepted during failures, but the acknowledgement came from nodes outside the key’s replica set, so R + W > N no longer describes the situation. It is availability bought by suspending the arithmetic, which is the right trade for a cart and the wrong one for a ledger.
When would you choose W=1, R=1?
For high-volume data where staleness is acceptable and write latency dominates the experience: telemetry, activity feeds, view counters, cache-like tables. You are explicitly accepting that a read may miss a recent write and that a node failure can lose an acknowledged write, so it must be data you can afford to lose or reconstruct.
How would you give a user read-your-writes on top of a quorum store?
Either read at a level that guarantees overlap for that key and rely on version resolution, or route that user’s reads to the replica or region that accepted the write for a short window, or have the client carry the version it wrote and retry until a replica reports at least that version. Session guarantees are a client-side construct; the quorum alone provides no per-session promise.
Answers that lose the round
- Stating R + W > N as if it delivered linearizability — it bounds staleness of completed writes and permits reading a value that later vanishes
- Omitting version metadata and expecting the quorum itself to identify the newest value
- Believing a failed quorum write leaves no trace; partial writes persist and can be resurrected by read repair
- Applying the arithmetic to a sloppy quorum, where the read and write sets may not intersect during a partition
- Choosing W=N for safety and creating a write path that any single slow or dead replica takes down
- Assuming last-write-wins is harmless — with concurrent updates and clock skew it silently discards data
- Forgetting anti-entropy, so replicas diverge permanently for keys nobody reads
Practise quorum reads and writes in a real repository
The quorum arithmetic exists to bound staleness, and this Gronex repository is where staleness bites: a write goes to the primary and the immediate follow-up read goes to a replica. The tests demand read-your-writes for the writer while keeping replica reads for everyone else, so the fix is a routing and versioning decision rather than a configuration flag.
FAQ
Is quorum the same as consensus?
No. Consensus uses quorums, but adds a protocol — a leader, terms, and a replicated log — that makes every replica agree on a single total order. A bare quorum read/write system has no agreed order, which is why it needs version reconciliation and why it cannot express compare-and-set without an extra consensus round.
What does `w: majority` mean in MongoDB?
That the write is acknowledged once a majority of the replica-set members holding data have applied it, which makes the write durable across a failover — a majority-acknowledged write cannot be rolled back, because any new primary must be elected by a majority that has seen it. Pairing it with `readConcern: majority` avoids reading data that a rollback could remove.
Does a bigger N make the system more consistent?
It makes it more durable, not more consistent. Consistency comes from the R and W choice relative to N. A larger N with quorum on both sides also raises the number of nodes each operation must contact, which raises tail latency, since you are waiting on the slowest of a larger sample.