Every Cassandra tutorial states the rule: if the number of replicas you read plus the number you write is greater than the replication factor, reads see the latest write. The rule is correct and almost useless on its own, because the questions engineers actually face are specific. What happens to a QUORUM write when one node is down? Why did a read return old data at LOCAL_QUORUM after a datacenter failover? Why did a write that timed out show up an hour later?

This article answers those questions by tracing eight scenarios replica by replica. Each one fixes a replication factor, a write level and a read level, injects a specific failure, and follows exactly which replicas hold which value and what the client sees. The consistency levels overview covers the vocabulary; this page is the workbook. It ends with code to set levels per query, a small simulator, a table for choosing levels by workload and a checklist.

Advertisement

The rule, and what it does not promise

A consistency level says how many replicas must respond before the coordinator answers the client. For a keyspace with replication factor RF, QUORUM means floor(RF / 2) + 1 replicas, summed across all datacenters; LOCAL_QUORUM means a quorum of the replicas in the coordinator's datacenter only. If W replicas acknowledged a write and R replicas answer a read, and R + W > RF, the two sets must share at least one replica, so the read sees the write. The coordinator resolves differences by the write timestamp of each cell: the highest timestamp wins.

Three things are not promised. Overlap does not order concurrent writers; two clients writing the same cell at once both succeed and the higher timestamp silently wins. A write that times out is not rolled back; it may be on some replicas and will spread. And the guarantee is about acknowledged writes reaching readers, not about linearizability, which needs lightweight transactions. Keep those three in mind; every surprising scenario below is one of them.

QUORUM write plus QUORUM read at RF=3: any two replicas out of three share at least one memberClientdriver, CL per queryCoordinatorany nodeReplica Av2 @ t=200Replica Bv2 @ t=200Replica Cv1 @ t=100 (missed)requestwrite ackwrite ackwrite lostWrite setA, BRead setB, COverlapB: returns v2The coordinator compares responses by write timestamp, returns v2, and (blocking read repair) writes v2 to C before replyingOverlap guarantees the newest acknowledged write is seen; it does not order concurrent writers or undo a failed write
Scenario 1 in one picture. The write reached A and B, the read asked B and C, and B is in both sets, so the newest value is returned and repaired onto C.

Scenario 1: QUORUM and QUORUM, RF=3, everything healthy

Replicas A, B and C hold v1 with timestamp 100. A client writes v2 at QUORUM with timestamp 200. The coordinator sends the mutation to all three replicas (it always sends to every live replica, whatever the level) and replies success after two acknowledgements. Suppose A and B ack quickly and C drops the mutation because it was briefly overloaded.

A read at QUORUM asks two replicas. One gets the full data and the other a digest. If it picks B and C, the digests differ, the coordinator fetches full data from both, sees timestamp 200 on B and 100 on C, returns v2 and, because the default table option read_repair = 'BLOCKING' is in force, writes v2 to C before replying. Once the client has seen v2 at QUORUM, no later QUORUM read can return v1. The Cassandra documentation calls this monotonic quorum reads.

Advertisement

Scenario 2: one replica down

Same keyspace, but C is down. A QUORUM write needs two live replicas; A and B are up, so it succeeds. The coordinator stores a hint for C and replays it when C returns, as long as C comes back within max_hint_window, which defaults to 3 hours. A QUORUM read also needs two and succeeds from A and B.

A write at ALL fails immediately with an Unavailable error: the coordinator knows from failure detection that only two of three replicas are alive and does not even try. That is different from a timeout, and the difference matters for retries. Unavailable means nothing was written. A timeout means the outcome is unknown.

RFLevelReplicas neededNodes that may be down
3ONE12
3QUORUM21
3ALL30
5QUORUM32
2QUORUM20

Scenario 3: ONE and ONE, and the stale read

A client writes v2 at ONE. A acknowledges; B and C both miss the mutation (the coordinator logged dropped mutations under load). Another client reads at ONE, and the driver's load-balancing policy happens to pick C. It returns v1. Nothing is broken: W + R = 2, which is not greater than 3, so no overlap was promised.

There is no read repair to rescue it either. Read repair only happens when the coordinator compares responses from several replicas, so it runs at levels such as QUORUM and LOCAL_QUORUM, not at ONE or LOCAL_ONE. The stale copy on C stays stale until a hint, a later read at a higher level or anti-entropy repair fixes it. ONE/ONE is the right choice for data where an occasional stale read is fine, such as metrics or a feed, and the wrong choice for anything a user just edited and is about to reload.

Scenario 4: the timed-out write that appears later

This is the scenario that causes the most confused bug reports. A client writes v2 at QUORUM. A applies it. B and C are slow, and after write_request_timeout (2,000 ms by default) the coordinator gives up and the client receives a write timeout reporting one acknowledgement out of two required.

The client assumes failure and tells the user the save did not work. A QUORUM read that hits B and C returns v1, which seems to confirm it. But v2 is sitting on A. The next QUORUM read that includes A returns v2 and blocking read repair copies it to the other replica, so the write the client thought had failed is now permanent. A hint or repair would have spread it anyway.

The rule that follows: a write timeout means unknown, not failed. Make writes idempotent, so retrying with the same value and the same client-supplied timestamp is harmless, and retry on timeout. If a write must either happen or not happen, that is a conditional update, which is scenario 7.

Scenario 5: two datacenters

Keyspace with NetworkTopologyStrategy, three replicas in dc1 and three in dc2, so the total RF is 6. QUORUM is floor(6 / 2) + 1 = 4. It is not two in each datacenter; four acknowledgements from anywhere will do, which means every QUORUM operation from dc1 must wait for at least one replica in dc2. Every request pays a cross-datacenter round trip, and if the link to dc2 fails, dc1 can still reach only three replicas and every QUORUM operation fails.

LevelReplicas needed (3 + 3)Cross-DC waitSurvives losing dc2
LOCAL_ONE1 localNoYes
LOCAL_QUORUM2 localNoYes
QUORUM4 anywhereYesNo
EACH_QUORUM (writes)2 in each DCYesNo
ALL6YesNo

The usual design is LOCAL_QUORUM for both reads and writes, with clients pinned to their local datacenter. That gives read-your-writes inside a datacenter, and only eventual consistency between them, typically within the cross-datacenter replication lag. The trap appears on failover: a user's write lands at LOCAL_QUORUM in dc1, dc1 fails a second later, their next request is routed to dc2, and dc2 may not have the write yet. If that matters, keep sessions sticky to one datacenter and fail over whole users rather than single requests. Multi-datacenter deployment covers the topology side.

Scenario 6: the RF=2 trap

A team chooses RF=2 to save disk and uses QUORUM for safety. QUORUM at RF=2 is floor(2 / 2) + 1 = 2, which is every replica. One node restart makes every token range it owns unavailable for QUORUM reads and writes, so a routine rolling restart becomes an outage. With RF=2 the realistic choices are ONE (available, no overlap) or ALL (overlap, no availability). Use RF=3 for anything that needs both, which is why it is the common default.

Scenario 7: lightweight transactions

Some writes must not be lost to a concurrent writer: claiming a username, decrementing stock, taking a lease. Timestamps cannot do that, so Cassandra offers conditional updates built on Paxos among the replicas. The condition is checked and the write applied as one linearizable step.

INSERT INTO users (username, user_id, created_at)
VALUES ('ada', 42, toTimestamp(now()))
IF NOT EXISTS;

UPDATE stock SET qty = 4 WHERE sku = 'X1' IF qty = 5;

Two levels apply. The serial consistency level, SERIAL or LOCAL_SERIAL, governs the Paxos phase: SERIAL involves a quorum of all replicas and LOCAL_SERIAL only the local datacenter's, with the same latency and failure trade-off as scenario 5. The regular consistency level governs the commit that makes the value visible to normal reads. Reading at SERIAL also finishes any in-progress Paxos round first, so it returns a value no other transaction can overturn.

The cost is several round trips between coordinator and replicas instead of one, and contention between clients on the same partition causes retries and timeouts; cas_contention_timeout defaults to 1,000 ms. Use conditional updates for the few rows that need them, never as a general locking layer. Recent Cassandra releases have reworked the Paxos implementation to need fewer round trips; check the documentation for your version before assuming the older costs. See lightweight transactions for the protocol.

Scenario 8: clock skew beats quorum

Two application servers update the same cell. Server S1's clock runs 300 ms fast. S1 writes v2 at real time 1,000 ms, stamped 1,300. S2 writes v3 at real time 1,100 ms, stamped 1,100. Both writes succeed at QUORUM, and every replica eventually holds both versions. Last-write-wins picks the highest timestamp, so v2 survives even though v3 was written later. No consistency level can fix this, because the replicas agree perfectly; they agree on the wrong winner.

The defences are to keep NTP or a better time source tight on every client, to avoid read-modify-write of a shared cell from several writers (model updates as new rows, for example appending events), and to use a conditional update where order truly matters.

Setting levels in code

Set a sensible default per application and override per statement. With the DataStax Python driver:

from cassandra import ConsistencyLevel, WriteTimeout, Unavailable
from cassandra.cluster import Cluster, ExecutionProfile, EXEC_PROFILE_DEFAULT
from cassandra.policies import DCAwareRoundRobinPolicy, TokenAwarePolicy
from cassandra.query import SimpleStatement

profile = ExecutionProfile(
    load_balancing_policy=TokenAwarePolicy(DCAwareRoundRobinPolicy(local_dc="dc1")),
    consistency_level=ConsistencyLevel.LOCAL_QUORUM,
    serial_consistency_level=ConsistencyLevel.LOCAL_SERIAL,
)
session = Cluster(["10.0.0.1"], execution_profiles={EXEC_PROFILE_DEFAULT: profile}).connect("shop")

# A cheap, staleness-tolerant read overrides the default.
feed = SimpleStatement("SELECT * FROM feed WHERE user_id = %s LIMIT 50",
                       consistency_level=ConsistencyLevel.LOCAL_ONE)

claim = SimpleStatement("INSERT INTO users (username, user_id) VALUES (%s, %s) IF NOT EXISTS")
try:
    result = session.execute(claim, ("ada", 42))
    print("claimed" if result.was_applied else "taken")
except Unavailable:
    print("not enough live replicas: nothing was written, safe to retry later")
except WriteTimeout:
    print("outcome unknown: read at SERIAL before deciding")

In cqlsh, CONSISTENCY LOCAL_QUORUM; and SERIAL CONSISTENCY LOCAL_SERIAL; set the levels for the session, which is useful for reproducing these scenarios by hand.

A simulator for overlap and availability

Before choosing levels for a new keyspace, enumerate every failure combination. This brute-force script checks, for a given RF, write level and read level, whether reads always overlap writes and how many replicas can be down while both still succeed.

from itertools import combinations

def needed(level, rf):
    return {"ONE": 1, "TWO": 2, "THREE": 3, "QUORUM": rf // 2 + 1, "ALL": rf}[level]

def analyse(rf, w_level, r_level):
    w, r = needed(w_level, rf), needed(r_level, rf)
    replicas = range(rf)
    overlap = all(set(ws) & set(rs)
                  for ws in combinations(replicas, w)
                  for rs in combinations(replicas, r))
    tolerated = rf - max(w, r)
    return overlap, tolerated

for rf, wl, rl in [(3, "QUORUM", "QUORUM"), (3, "ONE", "ONE"), (3, "ONE", "ALL"),
                   (2, "QUORUM", "QUORUM"), (5, "QUORUM", "QUORUM"), (5, "TWO", "QUORUM")]:
    ok, down = analyse(rf, wl, rl)
    print(f"RF={rf} W={wl:<6} R={rl:<6} overlap={ok!s:<5} tolerates {down} down")

Running it shows, for example, that RF=3 with ONE writes and ALL reads overlaps but tolerates no failed node for reads, and that RF=5 with TWO and QUORUM does not overlap (2 + 3 = 5, not more than 5).

Choosing levels by workload

WorkloadWriteReadWhy
User profile edits, single regionQUORUMQUORUMRead-your-writes, survives one node
Same, multi-region, sticky usersLOCAL_QUORUMLOCAL_QUORUMNo cross-region latency; cross-region is eventual
Metrics, logs, feedsLOCAL_ONELOCAL_ONEThroughput first, staleness acceptable
Unique claims, inventoryIF conditions, LOCAL_SERIALSERIAL or LOCAL_SERIALNeeds linearizability, not just overlap
Write once, read rarely, must not be missedONEALLCheap writes; reads pay, and fail if any replica is down

Operating it

  • Alert on client-side Unavailable and timeout counts per level; a rise in timeouts at QUORUM is the earliest sign of a slow replica.
  • Watch dropped mutations and pending hints. They tell you how much data is landing on fewer replicas than RF, which is what makes low-level reads stale.
  • Repair every table at least once per gc_grace_seconds window (ten days unless changed); otherwise a replica that missed a delete can resurrect the row. Hints cover only max_hint_window; see hinted handoff.
  • Do not change a level in production without rerunning the simulator for your RF and checking both overlap and the number of failures tolerated.

What to do next

  1. List each table, its RF per datacenter, and the read and write levels the application actually uses (grep the code, do not trust the design document).
  2. Run the simulator for each combination and flag any that do not overlap where the product expects read-your-writes.
  3. Make every write idempotent and treat write timeouts as unknown; add client-supplied timestamps where retries happen.
  4. Move multi-datacenter traffic to LOCAL_QUORUM with sticky routing, and decide explicitly what a user sees after a datacenter failover.
  5. Replace read-modify-write on shared cells with append-style rows or conditional updates, and check NTP on every client host.
  6. Reproduce scenarios 2, 3 and 4 on a test cluster with cqlsh and a stopped node, so the team has seen each one happen.
Key takeaway: Tunable consistency is overlap arithmetic plus three caveats: timestamps order writes, timeouts are not rollbacks, and overlap is not linearizability. Pick levels per workload with the arithmetic checked, prefer LOCAL_QUORUM with sticky routing across datacenters, avoid RF=2, use conditional updates only where order truly matters, and keep hints and repair healthy so low-level reads stay close to the truth.