Chapter 06 · Quorums, Consistency Levels, Read Repair, and Anti-Entropy
Quorum Edge Cases: Network Partitions, Stale Replicas, Failed Coordinators, and Timeouts
Trace quorum operations through partitions, slow replicas, coordinator failure, timeouts, ambiguous commit outcomes, and idempotent client retries.
Learning outcomes
Quorum requests fail in more ways than “success” or “database down.” A client can know that too few replicas are available, hit a deadline while work continues, lose the coordinator after replicas persisted the write, or retry an operation whose first outcome is unknown. This lesson gives AtlasMart a precise incident vocabulary for those histories.
Distinguish unavailable, timed out, known failed, successful, and successful-but-client-ambiguous outcomes.
Trace a quorum write through a network partition, slow replica, coordinator failure, and retry.
Explain why a client timeout does not roll back replicas that already persisted the operation.
Connect idempotency keys/request identities to safe retry behavior without claiming exactly-once execution.
Identify the logs, acknowledgement sets, versions, and coordinator/request IDs needed to reconstruct a quorum incident.
A timeout is information about the observer’s deadline, not proof that the remote operation failed. Chapter 04 added idempotency keys; Chapter 06 now applies that lesson to quorum writes where multiple replicas can commit before the response path fails.
1. Unavailable and timeout are not synonyms
Unavailable means the coordinator can determine that the requested rule cannot currently be satisfied—for example, fewer eligible/live replicas exist than the required threshold. A timeout means the deadline expired before enough qualifying evidence arrived. Some replicas may have completed their work. A client-ambiguous write is one where the caller cannot infer whether the operation took effect.
This distinction changes retry logic. Retrying a known non-executed operation can be straightforward. Retrying an ambiguous non-idempotent increment, charge, or shipment command can duplicate a side effect.
2. The coordinator is part of the response path, not necessarily the data owner
In leaderless systems, a coordinator may fan the request to replicas and collect acknowledgements. If the coordinator crashes after replicas persist the write but before the client receives success, the stored data does not automatically disappear. The client sees failure/timeout while replicas contain success.
A new coordinator handling the retry needs a stable request identity or a conditional/version rule to recognize that the logical operation already occurred. Otherwise “retry until success” can change the history.
3. Network partitions create asymmetric evidence
During a partition, one client may reach two replicas while another reaches a different subset. A strict quorum policy can reject one side; sloppy quorum may accept writes on fallback nodes; a low threshold may allow both sides to make progress and create concurrent versions. The correct interpretation depends on the product’s placement, version, and conflict semantics.
Operationally, record the routing/membership epoch, replica identities, request ID, version token, per-replica status, coordinator deadline, retry count, and whether hints/repair are pending. “Request failed” is too little evidence to diagnose the resulting state.
4. Deliberately wrong approach — retry a timed-out increment as a new command
AtlasMart increments a loyalty balance. A and B persist the increment, but the coordinator’s response is lost. The client times out and sends a new request with a new identity. A and B apply the increment again. Every quorum was locally satisfied, yet the business side effect happened twice.
The repair is to make the command idempotent or conditional at the business boundary: persist the request identity with the effect, reject duplicate identities, or model the operation as setting a versioned state rather than applying an unbounded increment. The retention scope of the deduplication record must cover the retry window.
5. AtlasMart lab — ambiguous commit and idempotent retry
from dataclasses import dataclass
@dataclass
class Request:
request_id: str
delta: int
counter = 10
seen = set()
replicas = {"A":10, "B":10, "C":10}
def apply(req, idempotent):
global counter
if idempotent and req.request_id in seen:
return "DEDUPLICATED"
# Coordinator sends to A and B; both persist before its reply is lost.
replicas["A"] += req.delta
replicas["B"] += req.delta
if idempotent:
seen.add(req.request_id)
return "COMMITTED_ON_A_B_BUT_REPLY_LOST"
req = Request("req-77", 1)
print("FIRST ATTEMPT", apply(req, idempotent=False), replicas)
print("client observes TIMEOUT: outcome is ambiguous, not known failure")
print("retry with NEW request identity ->")
req2 = Request("req-78", 1)
print(apply(req2, idempotent=False), replicas)
print("duplicate side effect on A/B?", replicas["A"] == 12)
print("\nREPAIRED POLICY")
replicas = {"A":10, "B":10, "C":10}; seen.clear()
req = Request("req-99", 1)
print("attempt 1", apply(req, idempotent=True), replicas)
print("same request id retry", apply(req, idempotent=True), replicas)
print("side effect count", replicas["A"] - 10)
print("\nSTATUS VOCABULARY")
print("UNAVAILABLE = not enough eligible/reachable replicas to meet the rule")
print("TIMEOUT = deadline expired before sufficient evidence arrived")
print("AMBIGUOUS = client cannot infer whether a timed-out write committed")
Verification checklist
- The first write updates A and B before the simulated reply is lost.
- The client-facing timeout does not revert A or B.
- A retry with a new request identity increments both replicas again.
- Reusing the same idempotency key in the repaired policy deduplicates the second attempt.
- The lab explicitly distinguishes unavailable, timeout, and ambiguous outcome.
Check your understanding
- Why can a timed-out write still be committed?
- What makes a retry dangerous?
- What does an idempotency key guarantee by itself?
- How is unavailable different from timeout?
- What incident evidence is essential?
Review the answers
1. The coordinator/client deadline may expire after replicas have persisted the operation but before enough acknowledgements reach the caller.
2. If the first outcome is ambiguous and the operation is not idempotent/conditional, the retry can apply the business effect again.
3. Nothing unless the service stores/checks it atomically with the effect for a defined retention scope; it is a protocol mechanism, not magic exactly-once delivery.
4. Unavailable is a determination that the threshold cannot be met from eligible/reachable replicas; timeout is a deadline expiration before sufficient evidence arrives.
5. Request/idempotency ID, coordinator, replica set, versions, per-replica acknowledgements, timeout point, retry history, membership epoch, and subsequent repair/hint state.
6. Production judgment
Retries need exponential backoff/jitter, bounded budgets, idempotency or conditional semantics, and backpressure so a partial outage does not become a retry storm. Observe timeout versus unavailable counts separately, plus coordinator saturation, queue depth, per-replica latency, duplicate-request detections, and unresolved version conflicts. Failure injection belongs in isolated test topologies, never on production replicas without a controlled game-day plan.
The final lesson turns these histories into measurements: how stale were reads, how long did replicas lag, and did any business invariant actually break?
Authoritative references
- Dynamo: Amazon’s Highly Available Key-value Store — coordinator, quorum-like reads/writes, sloppy placement, and version reconciliation
- Apache Cassandra 5.0 — Hints — current coordinator/replica and hinted-handoff example
- Apache Cassandra — Read repair — failed quorum write/read-repair edge case and monotonic quorum reads