skip to content

A replicated key-value store has N=5 replicas. A team configures writes to require acknowledgment from W=3 replicas and reads to query R=3 replicas. Why does this W+R>N configuration let the system mask node failures while still returning consistent reads, and what happens if the team instead sets W=2, R=2?

level: middleimportance: must knowfreq 62%

answer

  1. W+R>N guarantees quorum intersection
  2. pigeonhole: overlapping subsets can't both be disjoint if sizes sum > N
  3. majority quorum ties to 2f+1
  4. W=2,R=2 on N=5 -> stale reads possible
  5. read-repair reconciles overlap via version/vector clock

basics

~20 s

With W+R>N, any read set and any write set must share at least one replica, so a read always sees the latest write even if some nodes are down or slow. If W+R is not greater than N, like W=2,R=2 on N=5, reads and writes might miss each other entirely and return stale data.

solid answer

~50 s

Quorum systems mask individual node failures by requiring only a subset (W out of N, or R out of N) of replicas to respond rather than all N, so the operation succeeds despite some replicas being down or slow. The W+R>N condition guarantees any write quorum and any read quorum intersect in at least one replica - by pigeonhole, since W+R>N means they can't be fully disjoint within N nodes - so a read is guaranteed to overlap with the most recent successful write and can return (or reconcile via version comparison) the latest value. With W=3,R=3,N=5, overlap is guaranteed (3+3=6>5) and the system tolerates up to 2 failed replicas for either operation. With W=2,R=2, 2+2=4 is not greater than 5, so a read quorum can consist entirely of replicas that missed the last write, silently returning stale data - trading consistency for lower latency and higher availability.

go deeper

for a junior

Can state that quorum systems don't need every replica to respond, and can recall the W+R>N rule of thumb even if the reasoning isn't fully internalized.

for a middle

Can explain the pigeonhole argument for why W+R>N guarantees overlap, and can compute how many failures a given N/W/R combination tolerates.

for a senior

Can compare majority-quorum vs skewed-quorum configurations, connect the choice to CAP-theorem trade-offs, and explain why intersection alone doesn't give linearizability without versioning.

for a principal

Can reason about operational failure modes under partitions, how read-repair/hinted handoff restore consistency after the fact, and can justify a system's quorum defaults against specific latency/availability SLOs.

## What quorum masking is **Quorum-based fault masking** is a technique used in replicated systems (Dynamo, Cassandra, Riak, and conceptually classic Paxos/Raft majority quorums) to tolerate node failures without requiring every replica to be reachable for every operation. Instead of demanding unanimous agreement from all N replicas - which would mean a single down node blocks the whole system - the protocol only requires a quorum: - a subset of size `W` (**write quorum**) to acknowledge a write - a subset of size `R` (**read quorum**) to respond to a read As long as enough replicas are alive to form these subsets, the operation succeeds even though some replicas are unreachable, slow, or crashed; those replicas are effectively masked out from the client's perspective. ## Why the intersection property holds The correctness of this scheme for strong consistency hinges on the **quorum intersection property**: if `W+R>N`, then any write quorum and any read quorum, both drawn from the same pool of N replicas, are mathematically guaranteed to share at least one common replica. This falls directly out of the **pigeonhole principle**: two subsets of a size-N set whose sizes sum to more than N cannot be disjoint, because if they were disjoint their combined size (W+R) could be at most N. Because they intersect, at least one replica in the read set actually received and stored the most recent successful write, so the client can be given the latest value - either: - directly if the system tracks recency (via a version number or vector clock and picks the freshest response among the R replicas that answered) - by triggering **read-repair** to reconcile stale replicas ## Concretely, with N=5, W=3, R=3 Concretely, with `N=5`, `W=3`, `R=3`: `W+R=6>5`, so intersection is guaranteed, and the system tolerates: - up to `N-W=2` replicas being down for writes to still succeed - up to `N-R=2` replicas being down for reads to still succeed These can be different sets of 2 failed replicas across different operations, and the guarantee still holds because the math is about any quorum, not a fixed set of survivors. This is the essence of quorum-based fault masking: **redundancy absorbs failures transparently**, without the client or operator needing to detect which specific node is down - the protocol just needs enough live nodes to hit the W or R count, from any combination. ## What happens if the team instead sets W=2, R=2 If the team instead configures `W=2, R=2` on `N=5`, then `W+R=4`, which is not greater than `N=5`. Intersection is no longer guaranteed: it's possible to pick a write quorum of {A, B} and a read quorum of {C, D}, entirely disjoint, so the read never touches a replica that has the latest write and can return stale data with no way to detect it was stale. This isn't a bug so much as a deliberate trade-off: systems like Cassandra let operators choose weaker quorum configurations to lower latency (fewer replicas to wait for) and increase availability under partition or replica loss at the cost of losing the strong-consistency guarantee - a direct, concrete instance of the **availability/consistency trade-off**, expressed through the specific knob of quorum sizing. ## Subtleties beyond the inequality There are important subtleties beyond the basic `W+R>N` inequality. 1. First, quorum intersection guarantees you'll see the latest acknowledged write, but it does not by itself guarantee **linearizability** across concurrent operations - you also need a way to order or version writes so that when two write quorums partially overlap, the read can identify which value is actually newer rather than just returning some value that was written. 2. Second, **majority quorums** (the special case W=R=more than half of N) are the most common choice because they simultaneously minimize the replicas needed for both reads and writes while still guaranteeing intersection, and they compose naturally with N=2f+1 crash-fault tolerance - a majority write quorum can succeed as long as a majority of nodes are alive. 3. Third, real systems must handle the case where W or R momentarily can't be met, such as during a **network partition** that isolates a minority of replicas - the minority side simply cannot complete quorum operations, which is itself the fault-masking mechanism working as intended: it sacrifices availability on the minority side to preserve consistency, rather than allowing a stale or diverged partition to serve traffic silently.

  • Why is a majority quorum (more than half of N) the most common choice rather than picking arbitrary W and R that merely satisfy W+R>N?
    A majority quorum minimizes the maximum of W and R for a given fault tolerance while still guaranteeing intersection between any two majorities, and it directly maps to the classic n≥2f+1 crash-fault-tolerance bound used in Paxos/Raft. It also lets both reads and writes tolerate the same number of failures symmetrically, rather than skewing tolerance toward one operation type.
  • Does satisfying W+R>N by itself guarantee linearizable reads?
    No - it only guarantees that a read quorum overlaps with the most recent write quorum, so it will observe some replica holding the latest write. Without versioning such as timestamps or vector clocks, the read can't necessarily tell which of several returned values is actually the newest, so extra machinery is needed on top of plain quorum intersection for full linearizability.
  • During a network partition that isolates 2 of 5 replicas, with W=R=3, what happens to write availability on each side?
    The majority side (3 replicas) can still form quorums and continue serving both reads and writes normally. The minority side (2 replicas) cannot reach W=3 or R=3 on its own, so operations routed there fail or block rather than serving potentially stale or divergent data - the system chooses consistency over availability for that minority partition.

Think of five overlapping duty rosters guarding a warehouse: as long as any 'update the log' shift and any 'check the log' shift are big enough that they must share at least one guard, whoever checks the log will always run into someone who saw the latest entry. Shrink both shifts too much and you could have a checker who never crosses paths with the last person who updated it.

saying these in an interview costs you the question

  • Claims quorum masking means every replica must always respond
  • Thinks any W and R below N are automatically safe
  • Can't state the pigeonhole reasoning behind W+R>N
  • Believes quorum intersection alone guarantees full linearizability without versioning
  • Assumes lowering W and R has no consistency cost, only a performance benefit

context