Quorum Reads, Writes, and Tunable Consistency

LESSON

Consistency and Replication

006 30 min advanced

Quorum Reads, Writes, and Tunable Consistency

Harbor Point keeps the intraday trading limit for MUNI-77 on three replicas: iad, ord, and dub. At 09:30, the risk team changes the limit from 50M to 60M. One minute later, the reservation service must read that limit before accepting another order.

The simple rule sounds reassuring: “we use three replicas, so the new value is safe.” It is not enough. The write might have reached only two replicas, and the later read might ask only the third one. The database can truthfully say that replication is healthy while the reservation service sees the old limit.

For a single replicated key, quorum arithmetic makes the missing condition visible. It tells us whether every successful read must contact at least one member of every successful write. It does not make every replica current, decide which concurrent value is correct, or make a multi-key business rule atomic.

Core Insight

Let N be the number of replicas for one key, W the number of acknowledgements a write needs, and R the number of responses a read needs. When:

R + W > N

the successful read set and successful write set must share at least one replica. That is an overlap proof. It gives the reader a route to evidence of the write. It is not, by itself, a full consistency guarantee.

The operation-level choice matters. In Cassandra-like systems, a consistency level tells the coordinator how many replica responses to wait for; it is not a permanent adjective attached to the database. A safety-sensitive admission check can spend the latency of a quorum. A human dashboard may choose one fast replica and disclose that it can be stale.

The Small Model

For MUNI-77, Harbor Point uses N = 3. A quorum is a majority, so W = 2 and R = 2.

replicas:          iad   ord   dub
write needs W=2:    ✓     ✓
read needs R=2:     ✓           ✓
                    ^
                 overlap

The sets above both contain iad. More generally, two sets can avoid each other only if their total size is at most N. If R + W is larger than N, there is not enough room to keep them disjoint.

N W R R + W > N? What the arithmetic says
3 2 2 Yes Every successful read and write overlap.
3 2 1 No A one-replica read can choose the replica the write missed.
3 3 1 Yes Every read meets the all-replica write, but writes stop if one replica is unavailable.
5 3 3 Yes Two majorities overlap, with more failure tolerance than an all-replica write.

These numbers are a simplified model. They assume one fixed replica set for the key and a read that uses the same membership view as the write. Lesson 007 deliberately breaks that assumption with sloppy quorums and fallback nodes.

The Worked Path: Find the Overlap, Then Read It

Harbor Point writes the new limit at W=2.

starting state
  iad = 50M, v183
  ord = 50M, v183
  dub = 50M, v183

write: SET limit = 60M, version v184
  iad accepts v184
  ord accepts v184
  coordinator has W=2 and reports success
  dub is still at v183

Now compare two reads.

A read at one

read asks dub only
dub returns 50M, v183
coordinator returns 50M

This read did not fail. It followed its selected policy. But R=1 plus W=2 equals N=3, not greater than it. A legal read set {dub} and a legal write set {iad, ord} have no shared replica. The admission service could make a decision from an old limit.

A read at quorum

read asks ord and dub
ord returns 60M, v184
dub returns 50M, v183
coordinator compares versions
coordinator returns 60M, v184

Here the read set must include at least one member of the write set. That member exposes v184. The coordinator still has work to do: it must compare the returned versions and apply the store's resolution rule. It may also repair dub, but repair is a convergence action, not the reason the current read found 60M.

The evidence that corrects the initial model is concrete: the R=1 trace returned an allowed stale value. More copies were not the decisive fact. The selected R and W values determined whether the read was forced to meet the write.

What the Overlap Does Not Prove

The phrase “quorum read” can hide several further assumptions. Name them before turning the inequality into an API promise.

It does not choose the business-correct version

The overlap replica may contain a newer version according to the database's version rule. But “newer” needs a definition. Some systems use logical causality or preserve concurrent siblings. Cassandra documents a timestamp-based last-write-wins model for its mutations. A clock error or two concurrent updates can therefore produce a winning value that is valid under the storage rule but wrong for a business invariant.

For example, a risk analyst may concurrently lower a limit to 40M while an automated process raises it to 60M. Quorum overlap helps the reader encounter both replication histories. It does not decide whether a risk reduction should outrank an automated adjustment. That is a conflict-policy and authority question, developed later in the track.

It does not mean the write is durable under every failure

W=2 means the coordinator obtained two acknowledgements of the kind the database defines. Before claiming a survival guarantee, verify what those acknowledgements mean: memory receipt, local log durability, storage-engine commit, or something else. Lesson 005 showed why this boundary changes the failover story.

It does not make several keys one transaction

The inequality applies to the replica set for this key or partition. If an order approval reads a limit, creates a reservation, and appends an audit record in different partitions, each local quorum can overlap without making all three changes atomic. The service may need a transaction, one authoritative decision, escrow, or a compensating workflow. Do not extend a one-key proof into a cross-key promise.

It does not survive changed placement automatically

The proof expects a common home set. If a write uses fallback nodes during failure, or a read follows a different membership view, its candidate sets can fail to intersect even when the labels still say R=2, W=2, and N=3. That is the pressure for the next lesson on sloppy quorums and hinted handoff.

Choose a Level From the Contract

The familiar names ONE, QUORUM, and ALL are not a slider from cheap to good. They state which unavailable or stale-replica stories an operation accepts. Apache Cassandra, for example, defines QUORUM as a majority of the replicas and offers local and multi-datacenter variants whose replica sets differ. Vendor names must be interpreted together with replication factor, datacenter placement, and failure policy.

Harbor Point operation Possible policy Reasoning Boundary to state
Change an issuer limit Write W=2 of N=3 A later quorum safety check can overlap the acknowledged change. Confirm what write acknowledgement means and how conflicts resolve.
Check a limit before admission Read R=2 It must meet the acknowledged write set for that key. This alone does not atomically reserve against another concurrent check.
Refresh a trader dashboard Read R=1 Low latency may matter more than seeing the last update. Show a freshness label; do not reuse this read for admission.
Reconcile a stale replica Read enough replicas to compare versions, then repair as supported. The service needs evidence of divergence. Repair may be best effort and must not be assumed complete instantly.

This table is a situated recommendation, not a vendor recipe. A geographically distributed service might choose LOCAL_QUORUM to keep a client within one datacenter, accepting that the guarantee is scoped to local replication. A global business promise needs a separately specified cross-region write and read path.

Availability, Latency, and Failure

Stronger levels wait for more nodes. With N=3, a W=2 write can tolerate one unavailable replica. An ALL write cannot: one slow or unreachable replica prevents success. A R=2 read likewise needs two reachable responders, while R=1 can succeed with any one current member.

The trade-off is explicit: this buys overlap evidence, but costs more tail latency and reduces the set of failures through which the operation remains available. A good policy names the operation that deserves that cost; it does not upgrade every read because “quorum sounds safe.”

Watch the actual reasons a quorum request misses its contract:

Signal Interpretation
Requested versus achieved acknowledgements Whether the coordinator could form the intended W or R.
Per-replica response latency Which replica or network path consumes the quorum tail.
Replica version disagreement on reads Whether an overlapping read found divergence that needs repair.
Membership and topology changes Whether the assumed replica set still represents the key.
Timeout, unavailable, and fallback outcomes Whether the operation blocked, failed clearly, or used a weaker path.

Do not treat a timeout as proof that a write did not take effect. The coordinator may have received enough acknowledgements after the client stopped waiting. Commands that can be retried need an idempotency key or an outcome lookup, just as they did in the synchronous-replication case.

Check Your Understanding

Check: A key has N=5, W=3, and R=2. Does the arithmetic force every successful read to overlap every successful write?

Answer: No. R + W = 5, which is not greater than N. A write could use replicas A, B, C and a read could use D, E. Raising the read to R=3, or the write to W=4, creates the overlap condition.

Check: A W=2, R=2, N=3 design returns the version with the greatest timestamp. Does the equation prove that the returned limit is semantically correct after two concurrent updates?

Answer: No. It proves a shared replica exists between successful operations under the stated placement assumptions. The version and conflict rule decide which value wins; the business must decide whether that rule preserves the limit invariant.

Practice: Defend One Read Path

Harbor Point has N=3 replicas for an issuer limit. Its write path uses W=2. A product manager proposes using R=1 for both the trader dashboard and reservation admission because it is faster.

Write a short response that assigns a read policy to each endpoint. Include the arithmetic, the user consequence of a stale result, and one limit of the resulting guarantee.

A strong answer should say that the dashboard may use R=1 only if it displays or accepts staleness. The admission check should use R=2, since 2 + 2 > 3 forces overlap with a successful W=2 limit update. It should also state a boundary: the check still needs a sound version rule and does not make a concurrent multi-key reservation workflow atomic.

Connections

Resources

Key Takeaways

PREVIOUS Synchronous and Asynchronous Replication NEXT Leaderless Replication, Sloppy Quorums, and Hinted Handoff