Every Cassandra read and write carries a consistency level, and it is the most consequential number most applications set without thinking. It tells the coordinator how many replicas must answer before the request counts as a success. Get it right and you trade a little latency for exactly the guarantee you need; get it wrong and you either fail requests that could have succeeded or return data that is older than your users expect.
This page is the reference for the levels themselves: what each one means for reads and for writes, what the coordinator actually does while it waits, what the errors mean when it gives up, and how to set, enforce and observe levels in a real application. The quorum overlap argument and replica-by-replica scenarios have their own pages, linked at the end; here the goal is that you can pick a level for each query and know what will happen when a node is down.
What a consistency level controls
A consistency level is a property of a request, not of a table or keyspace. Replication factor, set per keyspace and per datacenter, decides how many copies exist. The consistency level decides how many of those copies the coordinator waits for. The two are independent: a write at ONE to a keyspace with replication factor 3 is still sent to all three replicas; the coordinator simply reports success after the first acknowledgement and lets the other two finish in the background.
That is the core idea to hold on to. Lower levels do not write less data; they wait for less evidence. A read at ONE asks one replica and believes it. A read at QUORUM asks a majority and returns the newest value among the answers, by write timestamp. Whether that newest value is the latest one ever written depends on whether the write and read sets overlap, which is the subject of the quorum rule.
Every level at a glance
The number the coordinator waits for is called blockFor. With replication factor RF in a single datacenter, a quorum is floor(RF / 2) + 1; with several datacenters, plain QUORUM uses the sum of every datacenter's RF, while the LOCAL_ and EACH_ variants compute a quorum per datacenter.
| Level | Writes | Reads | blockFor | Scope |
|---|---|---|---|---|
ANY | yes | no | 1, and a stored hint counts | any node |
ONE | yes | yes | 1 | any datacenter |
TWO | yes | yes | 2 | any datacenter |
THREE | yes | yes | 3 | any datacenter |
QUORUM | yes | yes | floor(sum of RF / 2) + 1 | all datacenters pooled |
LOCAL_ONE | yes | yes | 1 | coordinator's datacenter |
LOCAL_QUORUM | yes | yes | floor(local RF / 2) + 1 | coordinator's datacenter |
EACH_QUORUM | yes | yes, since 3.0 | a quorum in every datacenter | every datacenter |
ALL | yes | yes | every replica | all datacenters |
SERIAL | Paxos phase of conditional writes | serial reads | Paxos quorum | all datacenters |
LOCAL_SERIAL | Paxos phase of conditional writes | serial reads | Paxos quorum | coordinator's datacenter |
Two notes on that table. EACH_QUORUM reads were added in Cassandra 3.0 (CASSANDRA-9602); some older vendor documentation still lists it as write-only, so check the version a document describes before trusting it. And the serial levels are not ordinary consistency levels: on a conditional write they govern the Paxos round, while the normal consistency level governs the commit that follows.
As examples, with RF 3 in one datacenter QUORUM waits for 2. With RF 3 in each of two datacenters, QUORUM waits for 4 out of 6 from anywhere, LOCAL_QUORUM waits for 2 in the local datacenter, and EACH_QUORUM waits for 2 in each, 4 in total but with a stricter shape.
What the coordinator does while it waits
The node that receives the request becomes its coordinator. Before sending anything it checks, using the failure detector's view of which replicas are alive, whether enough replicas exist to possibly reach blockFor. If not, it fails immediately with an unavailable error and does no work. That check is cheap and it matters: an unavailable error means the request was never attempted.
Writes. The coordinator sends the mutation to every live replica of the partition. For each remote datacenter it sends one copy to a node there, which forwards it to the local replicas, so the write crosses the WAN once per datacenter rather than once per replica. It then waits until blockFor replicas, of the right datacenters, have acknowledged. Replicas that are known to be down get a hint stored on the coordinator for later delivery. If the timeout passes first, the client gets a write timeout, but the replicas that did receive the write keep it.
Reads. The coordinator picks blockFor replicas, preferring the fastest according to the dynamic snitch. It asks one of them for the full data and the others for a digest, a hash of their answer. If the digests match the data, it returns the result. If they differ, it fetches full data from those replicas, reconciles by timestamp, returns the newest values, and writes them back to the stale replicas it contacted, which is blocking read repair. A read waits for exactly blockFor answers; speculative retry can send an extra request if a replica is slow, but it does not raise the level.
# Simplified coordinator logic. Illustrative pseudocode, not Cassandra source.
def write(mutation, cl):
replicas = replicas_for(mutation.partition_key)
need = block_for(cl, replicas) # per datacenter for LOCAL_/EACH_
if live_count(replicas, cl) < need:
raise Unavailable(required=need, alive=live_count(replicas, cl))
for r in replicas:
send(r, mutation) if alive(r) else store_hint(r, mutation)
acks = wait_for_acks(need, timeout=write_request_timeout)
if acks < need:
raise WriteTimeout(received=acks, block_for=need) # some replicas may have applied it
return OKWhat the errors tell you
Most of the operational knowledge about consistency levels lives in the error types, because they tell you what state the data is in.
| Error | What happened | Data state | Safe to retry? |
|---|---|---|---|
| Unavailable | too few live replicas; nothing sent | unchanged | yes, ideally on another coordinator |
| WriteTimeout | sent, but fewer than blockFor acknowledged in time | unknown: applied on some replicas, maybe all | only if the write is idempotent |
| ReadTimeout | too few responses, or data replica silent | unchanged | yes |
| WriteFailure / ReadFailure | a replica answered with an error | depends on the error | investigate first |
The write timeout deserves emphasis because it is the source of most consistency surprises. A timeout is not a rollback. If two of three replicas applied the write and the third did not answer in time, the write exists, read repair and anti-entropy repair will spread it, and a later read will see it. Plain inserts and updates that set fixed values are idempotent and can be retried. Counter increments, list appends and non-idempotent logic are not. Timeout errors also report a write type such as SIMPLE, BATCH, BATCH_LOG, UNLOGGED_BATCH, COUNTER or CAS, which tells you which part of the path timed out.
Choosing a level
The rule that matters for freshness is that a read sees the latest acknowledged write when the read set and write set must overlap: R + W > RF, with R and W counted as replica numbers. In practice that gives a short list of sensible patterns.
| Pattern | Write | Read | Gives you | Costs you |
|---|---|---|---|---|
| Single datacenter, strong | QUORUM | QUORUM | reads see acknowledged writes; one replica of three may be down | latency of the second-fastest replica |
| Multi-datacenter, strong locally | LOCAL_QUORUM | LOCAL_QUORUM | the same, within each datacenter, without WAN latency | a reader in another datacenter may briefly see older data |
| Cross-datacenter durability | EACH_QUORUM | LOCAL_QUORUM | a write is acknowledged only once it is safe in every region | write fails if any datacenter is unreachable |
| High throughput, tolerant | LOCAL_ONE | LOCAL_ONE | lowest latency and best availability | stale reads after failures until repair |
| Write-heavy, read-strong | ONE | ALL | overlap without write latency | any single down replica fails reads |
Two warnings. Plain QUORUM in a multi-datacenter cluster pools every region's replicas, so with two datacenters most requests wait on a WAN round trip; it is rarely what you want. And with RF 2, a quorum is 2, so one node down makes every quorum request unavailable; use RF 3 for anything that needs quorum.
Setting levels in cqlsh and the Java driver
In cqlsh the level is a session setting, which makes it easy to test behaviour by hand:
cqlsh> CONSISTENCY LOCAL_QUORUM;
cqlsh> SERIAL CONSISTENCY LOCAL_SERIAL;
cqlsh> TRACING ON;
cqlsh> SELECT balance FROM bank.accounts WHERE id = 42;In the DataStax Java driver 4.x the defaults come from configuration. The shipped reference configuration sets basic.request.consistency to LOCAL_ONE, basic.request.serial-consistency to SERIAL and the request timeout to 2 seconds, so an application that never sets a level reads and writes at LOCAL_ONE. Set your own default and use execution profiles for exceptions:
# application.conf
datastax-java-driver {
basic.request.consistency = LOCAL_QUORUM
basic.request.serial-consistency = LOCAL_SERIAL
profiles {
analytics { basic.request.consistency = LOCAL_ONE }
}
}SimpleStatement debit = SimpleStatement
.builder("UPDATE bank.accounts SET balance = ? WHERE id = ?")
.addPositionalValues(newBalance, accountId)
.setConsistencyLevel(DefaultConsistencyLevel.LOCAL_QUORUM)
.setIdempotence(true) // fixed value: safe to retry
.build();
try {
session.execute(debit);
} catch (UnavailableException e) {
log.warn("only {} of {} replicas alive; not attempted", e.getAlive(), e.getRequired());
} catch (WriteTimeoutException e) {
log.warn("{} write: {} of {} acks; outcome unknown",
e.getWriteType(), e.getReceived(), e.getBlockFor());
}
ResultSet rows = session.execute(
SimpleStatement.newInstance("SELECT * FROM events.daily WHERE day = ?", day)
.setExecutionProfileName("analytics"));The driver's default retry policy is DefaultRetryPolicy. It also ships ConsistencyDowngradingRetryPolicy, which retries at a lower level; that quietly weakens the guarantee your code asked for, so use it only where you have decided stale results are acceptable.
Enforcing and observing levels
Since Cassandra 4.1 the server can police levels with guardrails (CASSANDRA-17188). In cassandra.yaml they are lists of levels to warn about or reject, empty by default:
# cassandra.yaml, 4.1 and later
read_consistency_levels_warned: [ONE, LOCAL_ONE]
read_consistency_levels_disallowed: [ANY]
write_consistency_levels_warned: []
write_consistency_levels_disallowed: [ANY, ALL]Warned levels succeed but return a client warning and a log line, which is a gentle way to find code paths still using a level you are phasing out. Disallowed levels are rejected. Guardrails apply to ordinary users, so roll out the warned list first, collect the offenders, then move levels to the disallowed list.
On the observation side, watch three things. Per-level timeouts and unavailables in the driver's metrics tell you which queries are failing and why. The server's client request metrics for reads and writes, with their Timeouts and Unavailables counters, tell you whether the cluster as a whole is struggling. And nodetool proxyhistograms shows coordinator-level latency, which is where a level that waits for a slow or distant replica shows up first.
Failure modes
- Assuming a failed write did nothing. A write timeout can leave the write visible later. Make writes idempotent, or read before retrying.
- Treating levels as isolation.
QUORUMreads plusQUORUMwrites do not make read-modify-write safe; two clients can both read 100 and both write 90. Use lightweight transactions for compare-and-set. - LOCAL_ONE reads after a failover. After a regional outage, a datacenter that missed writes serves old data at low levels until repair or read repair catches it up.
- ANY for important data. A write acknowledged only by a hint is invisible to reads until the hint is delivered, and lost if the coordinator holding it dies.
- Clock skew. Conflicts resolve by write timestamp, so a client or node with a fast clock can make an older write win regardless of level.
- Skipping repair. Low-level writes rely on hints, read repair and anti-entropy repair to converge. Without regular repair, deleted data can return once tombstones pass
gc_grace_seconds.
What to do next
- List every query path in your service with its current consistency level, including those relying on the driver default of
LOCAL_ONE. - For each path, write down whether it needs read-your-writes; where it does, confirm that R + W exceeds RF in the datacenter it runs in.
- Replace plain
QUORUMwithLOCAL_QUORUMin multi-datacenter clusters unless you have a stated reason for pooling regions. - Mark idempotent statements with
setIdempotence(true)on the builder (orsetIdempotent(true)on a built statement) and handle write timeouts explicitly for the rest. - Add the warned guardrail lists in a test environment, collect the warnings for a week, then decide what to disallow.
- Stop one replica in a staging cluster and run your workload at each level you use, checking which requests fail and whether any read returns stale data.
Continue with Cassandra consistency and quorum maths, Tunable Consistency Worked Examples, Cassandra Read Repair, Cassandra Hinted Handoff, Cassandra LWT architecture and Cassandra Speculative Retry.