Every introduction to leaderless replication gives the same rule: with N replicas, write to W and read from R, and if R + W > N every read overlaps the latest write. The rule is true and it is the starting point, covered with tunable consistency in quorum systems and with the theory of intersecting sets in Quorum Systems Explained. It is also where most misunderstandings start, because it describes sets of nodes, not what happens over time when requests overlap, nodes are slow and clocks disagree.
This article follows a quorum write and a quorum read through a real request path, using Cassandra-style terms with notes on Dynamo and Riak where they differ. It shows what the coordinator actually does, what an error really tells you, why quorum reads and writes are still not linearizable, how tail latency behaves, and how to pick levels for real workloads.
The write path, step by step
The client sends a write to any node, which becomes the coordinator for that request. Smart drivers route to a replica so the coordinator is also one of the owners. The coordinator hashes the partition key, looks up the N replicas in its token map, and sends the mutation to every live replica, not only W of them. W controls how many acknowledgements it waits for before answering the client, not how many replicas receive the data.
Each replica appends the mutation to its commit log, applies it to an in-memory table, and acknowledges. When W acknowledgements have arrived, the coordinator reports success. The remaining replicas still apply the write when their copy of the message arrives. If a replica is known to be down, or times out, the coordinator stores a hint and replays it when the replica returns, as described in hinted handoff.
Two details matter. First, every write carries a timestamp, usually assigned by the client or coordinator from its own clock, and replicas resolve conflicts by keeping the value with the highest timestamp: last write wins. Second, nothing coordinates the replicas with each other. There is no prepare phase and no commit decision. Each replica applies what it receives independently.
The read path: digests and read repair
For a read at level R, the coordinator asks the fastest-looking replica for the full data and asks R minus 1 others for a digest, a hash of their answer. If the digests match the data, the result goes straight back. If they disagree, the coordinator fetches full data from all R contacted replicas, merges them by timestamp to find the newest value, and sends the merged value back to any stale replica among them before replying to the client. This is blocking read repair, covered alongside background anti-entropy in read repair and anti-entropy.
Blocking matters. Without it, a read could return a new value from a replica that has it while the other replicas in a later read's quorum still hold the old one, so one client sees the new value and a second client, a moment later, sees the old one. Writing the newest value back to a quorum before returning means any later quorum read overlaps a replica that has it. Cassandra 4.0 made this behaviour configurable per table and removed the old probabilistic background read repair options; the default remains blocking.
A coordinator in code
The logic fits in a page. This sketch uses Python asyncio; replicas expose write and read coroutines that return timestamped values.
import asyncio
class Unavailable(Exception): pass
class Timeout(Exception): pass # outcome unknown, not failed
async def quorum_write(replicas, key, value, ts, w, timeout, hints):
live = [r for r in replicas if r.is_up()]
if len(live) < w:
raise Unavailable(f"{len(live)} live, need {w}") # nothing sent
for r in replicas:
if r not in live:
hints.add(r, key, value, ts) # known-down owner: hint now
tasks = {asyncio.create_task(r.write(key, value, ts)): r for r in live}
acks = 0
try:
for fut in asyncio.as_completed(tasks, timeout=timeout):
await fut
acks += 1
if acks >= w:
return "ok" # others keep running
except asyncio.TimeoutError:
pass
finally:
for t, r in tasks.items(): # store hints for laggards
t.add_done_callback(lambda f, r=r: f.exception() and hints.add(r, key, value, ts))
raise Timeout(f"{acks} of {w} acks; write may still be visible")
async def quorum_read(replicas, key, r, timeout):
chosen = replicas[:r] # real systems pick by latency
answers = await asyncio.wait_for(
asyncio.gather(*(x.read(key) for x in chosen)), timeout)
newest = max(answers, key=lambda a: a.ts)
stale = [x for x, a in zip(chosen, answers) if a.ts < newest.ts]
await asyncio.gather(*(x.write(key, newest.value, newest.ts) for x in stale))
return newest.value # repaired before returningNotice the two different errors. Unavailable is raised before anything is sent, so the write definitely did not happen. Timeout is raised after the mutation went out, so the write may have been applied anywhere from zero to N replicas. Cassandra drivers expose the same distinction as separate unavailable and write-timeout exceptions, and applications that treat both as failure are wrong about one of them.
What a timeout really means
Suppose W is 2 and only replica A acknowledges before the timeout. The client gets an error. But A has the new value, and B and C may apply it later from the original message or from a hint. Nothing rolls A back. A subsequent quorum read that contacts A and B finds A newer, repairs B, and returns the new value. The write that reported failure has become visible.
So a timed-out write is in an unknown state, and the only safe responses are to retry it or to read and check. Retrying is easy when the write is idempotent: setting a column to a value with the same timestamp twice is harmless. It is dangerous for counters and list appends, where a retry can apply the change twice. Design writes to be idempotent where possible, use client-assigned timestamps so a retry carries the same one, and keep non-idempotent operations away from paths that retry automatically.
Why quorums are not linearizable
Even with R + W > N and blocking read repair, a quorum store is not linearizable, for reasons the overlap rule does not cover.
- Concurrent writes lose data. Two clients write different values at nearly the same time. Replicas keep the higher timestamp, so one write silently disappears even though both clients saw success. Read-modify-write sequences, such as incrementing a balance, lose updates the same way.
- Clock skew reorders writes. Timestamps come from clocks. If a client whose clock runs 50 ms fast writes, and a second client writes 10 ms later in real time, the second write has the lower timestamp and loses, even though it happened after.
- Failed writes become visible later, as shown above, so a value can appear after the writer was told it failed.
- Sloppy quorums break the overlap. Dynamo and Riak can accept a write on substitute nodes outside the key's home replicas when owners are unreachable, counting them towards W. A later read from the home replicas may not overlap at all until the hints are delivered. Cassandra does not do this: hints never count towards the consistency level, except for the special
ANYwrite level, which can succeed with only a hint stored.
When you need a true compare-and-set, use a consensus protocol. Cassandra's lightweight transactions run Paxos per partition, with SERIAL or LOCAL_SERIAL consistency, at the cost of several extra round trips. Where concurrent updates must merge rather than overwrite, use vector clocks or CRDTs instead of last write wins.
Worked example: tail latency is an order statistic
A quorum operation finishes when the W-th fastest replica answers, so its latency is an order statistic of the replica latencies. Assume N is 3, replicas are independent, and each answers within 20 ms 99 percent of the time.
| Level | Slow when | Probability over 20 ms |
|---|---|---|
| ONE | all three replicas are slow | 0.01 cubed = 0.0001 percent |
| QUORUM (2 of 3) | at least two are slow | 3 x 0.01 squared x 0.99 + 0.01 cubed, about 0.03 percent |
| ALL | any one is slow | 1 - 0.99 cubed, about 3 percent |
Waiting for two of three makes the 20 ms tail about 30 times rarer than a single replica's, because one slow node is simply outvoted. Waiting for all three makes it three times more common: the operation's p99 is now worse than any replica's. This is why ALL is a poor default even before availability is considered. Real replicas are not independent, since garbage collection, compaction and shared networks correlate their slowness, so measured improvements are smaller, but the shape holds.
Reads have one more trick. Because the coordinator chose specific replicas, a slow one stalls the read. Speculative retry sends the request to an extra replica if the first ones have not answered within a threshold, such as the table's 99th percentile read latency, trading a little extra load for a much shorter tail.
Multiple datacenters
With replicas in several datacenters, plain QUORUM counts across all of them: with three replicas in each of two datacenters, it needs four of six, so every operation waits on the remote site. LOCAL_QUORUM counts only replicas in the coordinator's datacenter, two of three, and keeps latency local; data reaches the other site asynchronously. EACH_QUORUM for writes requires a quorum in every datacenter, which survives the loss of a site without losing acknowledged writes but makes writes as slow as the farthest site and unavailable if any one site is cut off.
Using LOCAL_QUORUM for both reads and writes gives read-your-writes within one datacenter, which is the right default for most applications whose users are pinned to a region. A user whose requests move between regions can read stale data until replication catches up, which is a product decision as much as an engineering one.
Failure modes and choices
| Symptom | Cause | What to do |
|---|---|---|
| Write reported failed but value appears | Timeout after partial apply | Treat timeouts as unknown; idempotent retries |
| Update lost under contention | Last write wins on concurrent writes | Lightweight transactions or mergeable types |
| Newer write lost to older one | Client clock skew | Synchronise clocks; server-side timestamps |
| Stale read after outage | Reads at ONE, or sloppy quorum before hints delivered | Quorum reads; run repair after outages |
| Deleted data returns | Replica down longer than the tombstone grace period, then repaired | Run full repair within the grace period |
| Unavailable errors with one node down | RF 2, or ALL levels | RF 3 per datacenter with quorum levels |
Hints are bounded in time, and a replica that was down longer than the hint window has missed writes no hint will deliver. Only anti-entropy repair, which compares Merkle trees between replicas, brings it back in line. Run repair on a schedule shorter than the tombstone grace period, and always after a long outage.
What to do next
- Set replication factor to 3 per datacenter and default reads and writes to LOCAL_QUORUM.
- Separate unavailable errors from timeouts in client code, and treat timeouts as unknown outcomes.
- Make writes idempotent with client-supplied timestamps, and keep counters and appends off automatic-retry paths.
- Use lightweight transactions only for genuine compare-and-set needs, and measure their latency.
- Monitor clock offset on every node and client host, and alert well below your tolerance.
- Enable speculative retry on latency-sensitive tables and schedule repair inside the tombstone grace period.