skip to content

questions

5

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

open as a page

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?

level: middleimportance: must knowfreq 70%

basics

~20 s

With 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.

open as a page

When tuning the read quorum size R and write quorum size W for a fixed replica count N in a quorum-replicated store, what are the concrete trade-offs between read latency, write latency, and consistency strength? Contrast R=1,W=N with R=N,W=1 and with a balanced choice like R=W=(N/2)+1.

level: seniorimportance: must knowfreq 65%

basics

~20 s

Making reads wait on more copies makes reads slower but writes faster, and vice versa for writes. Waiting on all copies for one side makes that side slow and fragile (one down node blocks it), while the balanced middle spreads the cost evenly across reads and writes.

open as a page

What is a 'sloppy quorum' in a quorum-replicated store like Amazon Dynamo or Cassandra, how does it differ from a strict quorum, and what consistency guarantee do you give up by using one?

level: middleimportance: should knowfreq 55%

basics

~20 s

A sloppy quorum lets the system write to any healthy nodes it can reach instead of insisting on the exact nodes normally responsible for that piece of data. It keeps writes working during outages, but it means you can no longer be sure a later read will overlap with that write.

open as a page

During a network partition that leaves fewer than R (or fewer than W) replicas reachable from a client's coordinator node, what happens to quorum reads and writes in a strict-quorum system, and how do techniques like sloppy quorums, hinted handoff, and read repair change that behavior? Also: name a case where R + W > N is satisfied yet a client can still observe stale or conflicting data.

level: principalimportance: should knowfreq 45%

basics

~20 s

If not enough copies are reachable, a strict system just refuses the read or write rather than risk giving a wrong answer — it picks correctness over always answering. Looser variants keep answering by writing to backup nodes and fixing things up later, but that reopens the door to occasionally seeing old or conflicting data.

open as a page