Cassandra does not have one consistency level. Every read and every write carries its own, chosen by the client, and the database does exactly what that level asks: wait for one replica, for a majority, for a majority in each data centre, or for all of them. That flexibility is why one cluster can serve a page-view counter and an account ledger, and it is also why consistency bugs in Cassandra are silent. No error is raised when an application writes at ONE and reads at ONE; it simply sometimes reads data that is older than what it just wrote.
This article explains the mechanics from first principles: what a coordinator does with a consistency level on the write and read paths, the overlap rule that turns levels into guarantees, how the levels behave across data centres, what timeouts really mean, how lightweight transactions differ, and which background processes converge replicas without guaranteeing anything. It finishes with a worked two-data-centre design, the failure modes worth rehearsing, and a checklist. It also corrects a claim often repeated: hints do count towards one level, ANY.
Replicas, coordinators and what a level means
Each partition key hashes to a token, and the replication strategy places the partition on a fixed set of replicas. With NetworkTopologyStrategy and {'dc1': 3, 'dc2': 3} every partition has three replicas in each data centre, placed on distinct racks where possible. The node that receives a request is the coordinator for it; with a token-aware driver that is usually one of the replicas itself.
A consistency level is a number of replica responses the coordinator waits for before answering the client. That is all it is. It does not change how many replicas receive a write: a write is always sent to every replica in every data centre. It changes only how many acknowledgements are required before the client hears 'success'. A replica acknowledges after appending to its commit log and applying to its memtable.
Replicas, coordinators and what a level means
Each partition key hashes to a token, and the replication strategy places the partition on a fixed set of replicas. With NetworkTopologyStrategy and {'dc1': 3, 'dc2': 3} every partition has three replicas in each data centre, placed on distinct racks where possible. The node that receives a request is the coordinator for it; with a token-aware driver that is usually one of the replicas itself.
A consistency level is a number of replica responses the coordinator waits for before answering the client. That is all it is. It does not change how many replicas receive a write: a write is always sent to every replica in every data centre. It changes only how many acknowledgements are required before the client hears 'success'. A replica acknowledges after appending to its commit log and applying to its memtable.
The levels
| Level | Replicas the coordinator waits for | Typical use |
|---|---|---|
ANY | Writes only. One replica, or a stored hint if none is up | Fire-and-forget telemetry; data may be unreadable for a while |
ONE, TWO, THREE | That many replicas, any data centre | Latency-sensitive, loss-tolerant data |
LOCAL_ONE | One replica in the coordinator's data centre | Reads that must never cross the WAN |
QUORUM | floor(total RF / 2) + 1 across all data centres | Single-DC strong reads and writes |
LOCAL_QUORUM | floor(local RF / 2) + 1 in the coordinator's DC | The default for multi-DC applications |
EACH_QUORUM | A quorum in every data centre | Writes that must be durable in every region before success |
ALL | Every replica | Rarely right: one down node fails the request |
SERIAL, LOCAL_SERIAL | Paxos quorum, for lightweight transactions | Compare-and-set, uniqueness |
Two details matter. QUORUM sums replication factors across data centres, so with 3 + 3 it needs 4 of 6 replicas and therefore at least one response from the remote region on every request. And ANY is the only level where a hint counts: the coordinator may return success having written the mutation only to its own hints, which means no replica can serve it until the hint is replayed.
The overlap rule
Guarantees come from arithmetic. If a write was acknowledged by W replicas and a later read consults R replicas out of RF, the two sets must share at least one replica whenever R + W > RF. That shared replica holds the new value, the coordinator reconciles responses by timestamp, and the read returns the newest acknowledged write. When R + W <= RF the sets can be disjoint and the read can legally return older data.
| Write CL (RF 3) | Read CL | R + W | Read sees latest acknowledged write? |
|---|---|---|---|
ONE | ONE | 2 | No guarantee |
ONE | ALL | 4 | Yes, but any down replica fails reads |
QUORUM | QUORUM | 4 | Yes, and tolerates one down replica on each side |
ALL | ONE | 4 | Yes, but any down replica fails writes |
LOCAL_QUORUM | LOCAL_QUORUM | 4 within one DC | Yes, for clients in the same DC |
The last row has a condition hidden in it. LOCAL_QUORUM overlap holds only within one data centre. A client in DC2 reading at LOCAL_QUORUM immediately after a client in DC1 wrote at LOCAL_QUORUM has no guarantee, because the write may still be crossing the WAN. Pin each user or entity to a home region, or use EACH_QUORUM writes when cross-region read-after-write matters.
The rule also assumes that timestamps order writes correctly. Cassandra resolves conflicts per cell by last-write-wins on the write timestamp, which the driver generates by default. Two clients with skewed clocks can make an older write win. Keep NTP tight, never set timestamps by hand without a reason, and use lightweight transactions when order really matters.
The read path and read repair
On a read at LOCAL_QUORUM with RF 3, the coordinator sends a full data request to the replica the dynamic snitch expects to be fastest, and digest requests (a hash of the result) to enough others to reach the level, here one. If the digests match, it returns the data. If not, it requests full data from the replicas it contacted, merges them cell by cell keeping the highest timestamp, writes the merged result back to the stale replicas it consulted, waits for those writes, and only then answers. This is blocking read repair, and it is what makes quorum reads monotonic: once a quorum read has returned a value, later quorum reads cannot return an older one. The table option read_repair = 'NONE' (Cassandra 4.0 and later) trades that property away for partition-level write atomicity; the old probabilistic read_repair_chance options were removed in 4.0. Details in read repair in depth.
Speculative retry hedges slow replicas: if the chosen replica has not answered within the table's speculative_retry threshold (the 99th percentile by default), the coordinator sends the request to another replica too. It improves tail latency; it does not change the guarantee.
Timeouts, unavailability and idempotence
Consistency levels produce two different errors, and conflating them causes real bugs. An UnavailableException is raised before anything is sent: the coordinator's failure detector already knows too few replicas are alive to satisfy the level. Nothing was written, and the request can be retried freely or degraded deliberately.
A WriteTimeoutException is raised after the mutation was sent and not enough replicas acknowledged in time. The write may have been applied on zero, some or all replicas, and Cassandra does not roll it back. A timeout is therefore an unknown outcome, not a failure. If the write is idempotent, for example setting a column to a value, retry it. If it is not, as with counter increments or list appends, a retry can double-apply it. Design writes to be idempotent, and mark statements idempotent in the driver only when they are.
from cassandra import ConsistencyLevel
from cassandra.cluster import Cluster, ExecutionProfile, EXEC_PROFILE_DEFAULT
from cassandra.policies import DCAwareRoundRobinPolicy, TokenAwarePolicy
profile = ExecutionProfile(
load_balancing_policy=TokenAwarePolicy(DCAwareRoundRobinPolicy(local_dc="dc1")),
consistency_level=ConsistencyLevel.LOCAL_QUORUM,
serial_consistency_level=ConsistencyLevel.LOCAL_SERIAL,
request_timeout=2.0,
)
session = Cluster(["10.0.0.1"], execution_profiles={EXEC_PROFILE_DEFAULT: profile}).connect("shop")
upsert = session.prepare("UPDATE carts SET items = ? WHERE cart_id = ?")
upsert.is_idempotent = True # a full overwrite is safe to retry after a timeout
session.execute(upsert, (items, cart_id))Avoid retry policies that silently downgrade the level after a failure. They turn a visible availability error into an invisible consistency violation, which is the worse of the two.
Lightweight transactions and serial consistency
Quorums give read-your-writes, not isolation. Two clients can both read 'username free' at QUORUM and both insert it. Lightweight transactions add a Paxos round so that a conditional write such as INSERT ... IF NOT EXISTS or UPDATE ... IF balance = 100 is linearizable for that partition. The serial consistency level (SERIAL or LOCAL_SERIAL) decides which replicas run Paxos; the ordinary consistency level decides how the committed value is written.
Linearizability costs extra round trips, and contention on one partition causes Paxos retries and timeouts. Cassandra 4.1 introduced an improved Paxos implementation that cuts round trips and fixed correctness issues; check which variant your cluster runs before relying on old latency numbers. Use LWTs for the rare decision that must be unique, not as a general transaction mechanism. See lightweight transactions.
Convergence is not consistency
Three mechanisms move replicas towards agreement; none of them is a consistency level. Hinted handoff: when a replica misses a write, the coordinator keeps a hint and replays it when the node returns, but only for outages shorter than max_hint_window (three hours by default). Read repair fixes what is read at levels above one. Anti-entropy repair compares Merkle trees of token ranges between replicas and streams the differences; it is the only mechanism that fixes data nobody reads after an outage longer than the hint window.
Repair also protects deletes. A tombstone is purged after gc_grace_seconds (ten days by default). If a replica missed the delete and repair does not run within that window, the old value survives on that replica and later spreads back as if the delete never happened. Every table therefore needs a full repair cycle shorter than its grace period. Read hinted handoff and repair for the operational side.
Worked example: one cluster, four levels
An e-commerce platform runs two regions, eu and us, each with RF 3. It has three kinds of data, and a single default level would be wrong for at least two of them.
Carts are read and written by one user who is routed to one home region. LOCAL_QUORUM for both gives read-your-writes for that user at LAN latency and survives one node down per region. If a region fails, users move to the other region and may see a cart a few seconds old, which the business accepts.
Orders must not be lost if a region disappears right after checkout. Writes go at EACH_QUORUM, so success means two replicas in each region have the order. The price is a WAN round trip on checkout and failed checkouts when either region lacks a quorum; the team decides that falling back to LOCAL_QUORUM with an alert is acceptable during a declared regional incident, and makes that an explicit runbook step rather than a driver retry policy.
Coupon redemptions must be unique. They use INSERT ... IF NOT EXISTS with LOCAL_SERIAL on a partition keyed by coupon id, and every coupon is homed in one region so two regions never race on the same Paxos state.
Product view counts are written at LOCAL_ONE and read at LOCAL_ONE; a stale or lost increment is cheaper than doubled latency. Each choice is written down next to the table definition so the next engineer knows it was deliberate.
Failure modes
- Writes at ONE, reads at ONE. Data appears and disappears between requests. Detect it by auditing driver profiles, not by waiting for user reports.
- Global QUORUM in a multi-DC cluster. Every request pays WAN latency and a regional partition makes the application unavailable everywhere.
- Retrying non-idempotent writes on timeout. Counters drift upward and list columns gain duplicates.
- Outages longer than the hint window with no repair. Replicas stay divergent indefinitely; deleted data resurrects after
gc_grace_seconds. - Clock skew. A client with a fast clock makes its writes unoverwritable until real time catches up.
- Cross-region read-after-write at LOCAL levels. A user who switches region sees an older state. Route users to a home region.
Trade-offs
Every level trades latency and availability against the strength of what a read can see. Higher write levels make reads cheaper to keep consistent and lower write availability; higher read levels do the reverse. The usual sweet spot in one region is LOCAL_QUORUM on both sides, because it tolerates a node loss per replica set without giving up overlap. Moving up to EACH_QUORUM buys regional durability at WAN latency; moving down to ONE buys latency at the cost of guarantees. More examples are in multi-DC Cassandra.
What to do next
- List every driver execution profile and statement-level override and record the read and write level for each table.
- Check R + W > RF for every table where correctness depends on reading the latest write.
- Replace global
QUORUMwithLOCAL_QUORUMorEACH_QUORUMin multi-DC clusters. - Remove downgrading retry policies and mark only truly idempotent statements as idempotent.
- Confirm that a full repair completes within
gc_grace_secondsfor every table. - Alert on hints stored and on nodes down longer than
max_hint_window. - Move uniqueness checks to lightweight transactions and keep them off hot partitions.
- Rehearse a node loss and a region loss in staging and observe which requests fail and which go stale.