skip to content

In a replicated data store, what is a read/write quorum, and what does the condition R + W > N guarantee (where N is the number of replicas, W is the write quorum size, and R is the read quorum size)?

level: juniorimportance: must knowfreq 80%

answer

  1. pigeonhole overlap
  2. R+W>N
  3. W = write acks needed
  4. R = read replicas queried
  5. N=3,R=2,W=2 common default

basics

~20 s

A quorum is the minimum number of copies of your data that must respond before a read or write counts as 'done'. If the number of copies you write to plus the number you read from is bigger than the total number of copies, you're guaranteed to always read at least one copy that has the newest write.

solid answer

~50 s

With N replicas holding copies of each record, a write doesn't wait for all N to ack — it waits for W acks. A read doesn't query all N — it queries R replicas and picks the newest version among the responses (via a version number, timestamp, or vector clock). The rule R + W > N guarantees the R replicas you read and the W replicas you wrote to must share at least one common replica, by the pigeonhole principle: two subsets of an N-element set whose sizes sum to more than N cannot be disjoint. That shared replica is guaranteed to hold the latest acknowledged write, so the read set as a whole always contains the newest version, letting the client detect and return it. It's a tunable point between single-copy strong consistency (R=N or W=N) and loose eventual consistency (R=1, W=1).

go deeper

for a junior

Should state that quorum means 'not all replicas needed' and that R+W>N is the rule for guaranteeing a read sees the latest write, even without deriving why.

for a middle

Should be able to explain the pigeonhole/overlap argument concretely with a worked N/R/W example, and know that reads reconcile versions rather than trusting a single replica.

for a senior

Should reason about R and W as independent latency/availability knobs, cite a realistic default (e.g. N=3,R=2,W=2), and know what happens when a quorum can't be reached.

for a principal

Should place quorum consistency within the CAP/PACELC trade space, discuss its limits (doesn't give linearizability outright), and know how production systems layer read repair, hinted handoff, and anti-entropy on top of it.

## The two numbers A quorum system sits underneath most 'eventually consistent but tunable' replicated stores — Dynamo, Cassandra, Riak, and their descendants. Each logical record is copied onto N physical replicas (nodes). Instead of requiring every replica to participate in every operation, the system defines two smaller numbers: - **W, the write quorum.** A write is considered successful once W replicas have durably stored the new value and acknowledged it back to the coordinator; the client doesn't wait for the remaining N-W replicas, which receive the update asynchronously (or via background repair). - **R, the read quorum.** A read queries R replicas, collects their versions of the value, and — because different replicas can be at different points in their update history — reconciles the responses by picking whichever version is 'newest' according to some ordering mechanism: - a monotonically increasing version counter - a wall-clock timestamp (risky under clock skew) - a vector clock that can detect true concurrency versus a simple ordering ## Why the overlap is guaranteed The reason R + W > N matters is a direct application of the **pigeonhole principle**. Imagine the N replicas as N boxes: 1. A write touches W of them; a read touches R of them. 2. If R + W > N, then the two groups cannot be chosen disjointly — there literally aren't enough boxes left to keep them apart — so at least one replica must appear in both the write group and the read group. 3. That shared replica received the newest write (it was one of the W that acknowledged it) and it is also one of the R replicas the read consults, so its version is present among the values the read collects. 4. As long as the client correctly identifies 'newest' among what it received, it will surface that value. This is why R + W > N is often phrased as 'the read quorum and write quorum must overlap.' ## Why not full replication Why build systems this way instead of just replicating synchronously to every node (N-of-N) or accepting a single writer? - **Full replication (W = N)** gives strong consistency but means a single slow or down node stalls every write — poor availability and poor tail latency. - **A single copy (N=1)** is fast and simple but has no fault tolerance at all. Quorums let an operator dial in a position between those extremes: N is chosen for durability and fault tolerance (how many node failures you can absorb before losing data), while R and W are chosen per-operation-type to trade off latency, throughput, and consistency strength independently of N. A write only needs a majority-ish subset to succeed, so one lagging or crashed replica doesn't block the whole system, and a read likewise only needs enough replicas to guarantee it sees the newest data, not all of them. ## The trade-off in the numbers The trade-off is visible in the numbers themselves. - **Larger R or W** means more replicas have to respond before the client gets an answer, which raises latency (bounded by the slowest of the R or W replicas contacted, so tail latency worsens as the quorum size grows) and lowers availability (more replicas must be reachable, so the system tolerates fewer simultaneous failures before an operation can't be completed at all). - **Smaller R or W** is faster and more available but weakens the consistency guarantee — if R + W <= N, quorums can be chosen so they never overlap, and a read can return stale data even though a newer write already succeeded elsewhere. A very common balanced choice is **N=3 with R=W=2** (used as the default 'QUORUM' consistency level in systems like Cassandra), which tolerates one replica being down for either reads or writes while still guaranteeing R+W(=4) > N(=3). ## Failure modes at the edges Failure modes show up at the edges. If a network partition or node crash drops the number of reachable replicas below R (for a read) or below W (for a write), a strict-quorum implementation simply fails that operation — it will not proceed with fewer acks than the configured quorum, because doing so would break the overlap guarantee. That is a deliberate **CAP-theorem trade**: quorum systems configured with R+W>N lean toward consistency over availability during a partition. Some systems relax this with **'sloppy quorums'** that write to whatever healthy replicas are reachable and reconcile later — trading the R+W>N guarantee for higher availability. A concrete example: Amazon's Dynamo paper (2007) popularized N=3, R=2, W=2 for a shopping-cart-style key-value store, explicitly choosing that balance so a single node failure never blocks reads or writes, while still making stale reads rare in the common case.

  • If N=5, W=3, R=2, is the read guaranteed to see the latest write?
    No — R+W = 5 = N, not greater than N, so the write set and read set can be chosen to exactly partition the 5 replicas with zero overlap. You'd need R=3 (or W=4) to force overlap and guarantee freshness.
  • Does R + W > N alone give you linearizability?
    Not fully. It guarantees the read set overlaps the last acknowledged write's replica set, but it doesn't order concurrent writes, doesn't handle a coordinator crashing mid-write before reaching W replicas, and doesn't prevent a read from racing an in-flight write. Real systems add version vectors, read repair, and sometimes a consensus layer for stronger guarantees.
  • Why not just always set R=1 and W=N, or R=N and W=1?
    Both extremes still satisfy R+W>N and guarantee overlap, but they push all the latency and availability cost onto one side of the operation — R=1,W=N makes every write wait for every replica (slow, fragile writes, fast reads), while R=N,W=1 makes every read wait for every replica (fast writes, fragile reads). A balanced quorum spreads that cost.

Think of N as the total number of people on a group chat who each got a copy of a memo. A 'write' is complete once W of them confirm they've read the latest version; a 'read' means you ask R of them what the memo says. If W+R is bigger than the group size, whoever you ask during your read must include at least one person who was in the group that confirmed the latest memo — there just aren't enough people left to dodge them.

saying these in an interview costs you the question

  • Says R+W>N means all replicas are always consistent immediately
  • Doesn't know W and R are independently tunable per operation type
  • Thinks quorum size must equal N/2 rounded up with no other valid choices
  • Can't explain the pigeonhole/overlap reasoning when asked why R+W>N specifically
  • Confuses quorum reads with leader-based replication (thinks there must be a single leader)

context