The usual summary says Cassandra is an AP system. That is true of its defaults and misleading as a guide to operating it. Cassandra makes the consistency-versus-availability choice per request, through the consistency level the client sends, and the right choice depends on replication factor, datacentre layout and exactly which failure is happening.
This article is a practical playbook. It works through the quorum arithmetic, what the coordinator does when replicas are missing, what node, rack and WAN failures do to each consistency level, how replicas converge afterwards, when lightweight transactions are worth their cost, and how client retries interact with all of it. The general introduction to CAP and Cassandra is in the CAP theorem article.
CAP is decided per request
CAP says that during a network partition a system must either refuse some requests (choose consistency) or answer them with possibly stale or divergent data (choose availability). In Cassandra, each read and write names a consistency level: how many replicas must acknowledge before the coordinator reports success. A low level stays available through more failures; a high level refuses requests sooner but gives stronger guarantees. The same table can be written at ONE by a metrics pipeline and read at LOCAL_QUORUM by a user-facing service.
Most of the time there is no partition, and the real trade-off is the one PACELC names: else, latency versus consistency. Waiting for two replicas instead of one adds the latency of the slower of the two to every request. So a consistency level is a choice about everyday latency as much as about failure behaviour.
The arithmetic
A quorum of N replicas is floor(N / 2) + 1. With a replication factor of 3, a quorum is 2, so quorum operations survive one replica down. Reads and writes are guaranteed to overlap in at least one replica when the read count plus the write count exceeds the replication factor. That is why QUORUM writes with QUORUM reads give read-your-writes for a single client, and why ONE with ONE does not.
| Level | Replicas needed (RF 3, one DC) | Replicas needed (RF 3 + 3, two DCs) | Survives |
|---|---|---|---|
| ONE / LOCAL_ONE | 1 | 1 (LOCAL_ONE: 1 in local DC) | All but one replica down |
| QUORUM | 2 | 4 of 6, any DCs | One DC down only if the other has 4+; a 3/3 split fails |
| LOCAL_QUORUM | 2 | 2 of 3 in the coordinator's DC | One local replica down; any remote failure |
| EACH_QUORUM | 2 | 2 in each DC | One replica per DC; not a WAN partition |
| ALL | 3 | 6 | Nothing |
| SERIAL / LOCAL_SERIAL | 2 | 4 of 6 / 2 of 3 local | As QUORUM / LOCAL_QUORUM |
Note the trap in the second column: with 3 replicas in each of 2 datacentres, QUORUM needs 4 of 6, which neither side of a WAN partition has. Multi-datacentre applications almost always want LOCAL_QUORUM. EACH_QUORUM is mainly used for writes that must land in every datacentre, and for reads or writes it fails during a WAN split.
What the coordinator does when replicas are missing
The node that receives a request becomes its coordinator. Before sending anything, it checks its failure detector: if it already believes fewer replicas are alive than the level needs, it fails fast with an UnavailableException. Nothing was written, so the request can be retried safely later.
If enough replicas look alive, the coordinator sends the write to all replicas for the row and waits for the required number of acknowledgements, up to the write timeout (2 seconds by default; reads default to 5). If not enough answer in time, the client gets a WriteTimeoutException. This is the case people misread. The write may have been applied on some replicas, and Cassandra does not roll it back; those replicas will spread it to the others through the mechanisms below. A failed write is an unknown outcome, not a rolled-back one.
Scenario playbook
Walk through the common failures for the layout in the diagram: two datacentres, three replicas each, racks mapped so that each rack holds one replica of every row, and clients using LOCAL_QUORUM.
- One node down. Every row it replicates still has 2 of 3 local replicas, so
LOCAL_QUORUMsucceeds. Coordinators store hints for the missing node. Nothing visible to users. - One rack down. With the rack-aware placement described in multi-datacentre replication, the rack held one replica of every row, so this is the same as one node down, for all rows at once. Without rack awareness, some rows could lose two replicas and
LOCAL_QUORUMwould fail for them. - Two local replicas of a row down.
LOCAL_QUORUMfails for that row with Unavailable. You can fall back toLOCAL_ONEfor reads that tolerate staleness, or route the request to the other datacentre. Doing either automatically is a product decision: it trades correctness for uptime. - WAN partition. Both sides keep serving
LOCAL_QUORUMtraffic and accept conflicting writes to the same rows. When the link heals, each cell resolves by last write wins on its timestamp. Anything that needs a single global order, such as unique usernames or stock counts, must not rely on ordinary writes during a partition. - Whole datacentre lost. Clients in that datacentre fail; clients elsewhere keep working. Driver failover to a remote datacentre is possible but, for local levels, means the remote side's view is authoritative, which is the same trade.
How replicas converge afterwards
Availability during a failure is paid for with divergence that must be repaired. Three mechanisms do it. Hinted handoff: the coordinator keeps writes for a replica that is down and replays them when it returns, but only for the hint window, 3 hours by default. Read repair: when a read at a level above ONE finds replicas disagreeing, the coordinator writes the newest value back to the stale replicas before answering; since Cassandra 4.0 this is blocking and there is no longer a random background read repair chance. Anti-entropy repair: nodetool repair compares replicas with Merkle trees and streams differences, and it is the only mechanism that reaches data nobody reads.
Repair has a deadline. Deletes are written as tombstones that are kept for gc_grace_seconds, 10 days by default. If a replica misses a delete and is not repaired before the tombstone is purged elsewhere, the old value comes back: a zombie. So every table needs a full repair cycle shorter than its grace period, and a node that has been down longer than the hint window needs a repair before or immediately after it rejoins. The mechanics are covered in hinted handoff and read repair.
When you need linearizability: lightweight transactions
Some operations cannot tolerate last write wins: creating an account with a unique name, claiming a job, reserving the last item. Cassandra's lightweight transactions (IF NOT EXISTS, IF column = value) run a Paxos round among the row's replicas, making the operation linearizable for that partition. This is choosing consistency on purpose. The serial level decides the scope: SERIAL needs a quorum of all replicas across datacentres and therefore fails on both sides of a WAN partition; LOCAL_SERIAL needs only the local datacentre and stays available, but is linearizable only among clients that also use it in that datacentre.
The cost is several round trips instead of one and contention when many clients hit the same partition. Cassandra 4.1 added a newer Paxos implementation, selected in configuration, that reduces round trips; general multi-partition transactions (the Accord work) are not part of Cassandra 5.0, so check your version before designing around them. Use lightweight transactions for the few rows where uniqueness or compare-and-set is essential, never as a default. Detail is in lightweight transactions.
-- Keyspace with three replicas in each datacentre.
CREATE KEYSPACE shop WITH replication = {
'class': 'NetworkTopologyStrategy', 'east': 3, 'west': 3
};
-- cqlsh: choose levels per session to experiment.
CONSISTENCY LOCAL_QUORUM;
SERIAL CONSISTENCY LOCAL_SERIAL;
-- Ordinary write: last write wins by timestamp, no read-before-write.
UPDATE shop.cart SET qty = 3 WHERE cart_id = 42 AND sku = 'A12';
-- Lightweight transaction: Paxos round among replicas, returns [applied].
INSERT INTO shop.orders (order_id, cart_id, status) VALUES (9001, 42, 'NEW') IF NOT EXISTS;
Clients: retries, idempotence and timestamps
Because a timeout leaves the outcome unknown, the client's retry behaviour is part of your consistency model. In the DataStax Java driver 4, statements are treated as non-idempotent unless you mark them otherwise, and the driver will not retry non-idempotent statements after write timeouts or run speculative executions for them. Mark a statement idempotent only when running it twice has the same effect as once: setting a column to a value is idempotent; incrementing a counter, appending to a list and most lightweight transactions are not.
// DataStax Java driver 4.x. application.conf:
// datastax-java-driver.basic.load-balancing-policy.local-datacenter = east
// datastax-java-driver.basic.request.consistency = LOCAL_QUORUM
// datastax-java-driver.basic.request.serial-consistency = LOCAL_SERIAL
// (basic.request.default-idempotence is false unless you change it)
SimpleStatement setQty = SimpleStatement
.builder("UPDATE shop.cart SET qty = ? WHERE cart_id = ? AND sku = ?")
.addPositionalValues(3, 42L, "A12")
.setIdempotent(true) // same values on retry: safe to retry and speculate
.build();
SimpleStatement placeOrder = SimpleStatement
.builder("INSERT INTO shop.orders (order_id, cart_id, status) VALUES (?, ?, 'NEW') IF NOT EXISTS")
.addPositionalValues(9001L, 42L)
.build(); // not idempotent by default: the driver will not blindly retry
try {
session.execute(setQty);
ResultSet rs = session.execute(placeOrder);
if (!rs.wasApplied()) { /* order already exists: read it and continue */ }
} catch (UnavailableException e) {
// Coordinator knew too few replicas were alive. Nothing was written. Safe to retry later.
} catch (WriteTimeoutException e) {
// Some replicas may have applied the write. Outcome unknown: re-read or retry only if idempotent.
}Last write wins uses the write timestamp, which by default comes from the client driver or the coordinator's clock. Clock skew between application hosts can make an older write win. Keep NTP tight, avoid having two writers race on the same cell, and model conflicting updates as separate rows (an event per change) where the order matters.
Worked example: a two-datacentre shop
A shop runs in east and west with RF 3 in each, clients pinned to their local datacentre. Carts are written and read at LOCAL_QUORUM with idempotent statements. Order placement uses IF NOT EXISTS at LOCAL_SERIAL, with order IDs allocated per datacentre from disjoint ranges so the two sides can never collide. Stock is not decremented in Cassandra during checkout; a reservation service that owns stock per warehouse does that.
During a 20-minute WAN partition, both datacentres keep taking orders. A user who edits a cart in east and then, after a DNS change, in west, may see an older cart: an accepted cost. After the link heals, hints from the 20 minutes replay well within the 3 hour window, and read repair fixes carts as they are read. The weekly repair schedule keeps every table under its 10-day grace period. Latency stays at one local round trip plus the slower of two local replicas, typically a few milliseconds, because nothing waits for the WAN.
Failure modes
| Symptom | Cause | Fix |
|---|---|---|
| Outage in both DCs during a WAN cut | QUORUM or SERIAL with replicas split evenly | LOCAL_QUORUM and LOCAL_SERIAL |
| Stale reads after a successful write | Read plus write replicas not exceeding RF, e.g. ONE/ONE | Quorum on both sides, or accept it explicitly |
| Deleted rows reappear | Repair cycle longer than gc_grace_seconds | Schedule repair within grace; repair returning nodes |
| Duplicate side effects | Non-idempotent statement retried after a timeout | Idempotent design; mark only safe statements |
| Newer update lost | Clock skew between writers | Tight NTP; single writer per cell; event rows |
| LWT latency spikes | Contention on hot partitions | Spread keys; limit LWT to essential rows |
The consistency levels themselves, including their read and write paths, are covered in Cassandra consistency levels.
What to do next
- Write down the replication factor per datacentre for each keyspace and compute the replicas each level needs, including under a WAN split.
- Move multi-datacentre applications from QUORUM to LOCAL_QUORUM and from SERIAL to LOCAL_SERIAL unless you need a global order.
- Treat WriteTimeoutException as an unknown outcome; mark only truly idempotent statements as idempotent.
- Use lightweight transactions only for uniqueness and compare-and-set, and design keys to avoid contention.
- Schedule repair so every table completes a cycle within gc_grace_seconds, and repair nodes that were down longer than the hint window.
- Test a node, rack and WAN failure in staging and record which requests fail and how long convergence takes.