Chapter 09 · The Write Path: Coordinator, Commit Log, Memtables, Replicas, and Acknowledgment
Client Request to Coordinator: Token Awareness, Replica Selection, and Write Consistency
Trace an AtlasMart write from token-aware client routing through a request-scoped coordinator to replicas, then prove consistency-level acknowledgment behavior under controlled replica loss.
Learning outcomes
AtlasMart receives an order-status mutation from an application
instance in dc1. The application needs to know
which node it contacted, which replicas own the partition, how
many acknowledgments are required, and why a successful response
is not proof that every replica is already identical. This
lesson separates routing, coordination,
replication, and consistency so “Cassandra
accepted my write” becomes an evidence-backed statement.
Explain the roles of client driver, coordinator, natural replica, token, and consistency level without treating them as synonyms.
Relate a partition key to token ownership and the replica set selected by NetworkTopologyStrategy.
Calculate the acknowledgments required by ONE, LOCAL_QUORUM, QUORUM, and ALL for RF=3 in one datacenter.
Use cqlsh tracing and nodetool evidence to distinguish coordinator work from replica work.
Run a reversible one-replica failure experiment and state what a successful quorum write does and does not guarantee.
The mandatory labs use the pinned
cassandra:5.0.9 image and the Java 17,
cqlsh, and nodetool versions bundled
by that image. The shared course cluster is
atlasmart-course with three disposable nodes
(atlasmart-cass-1..3) on Docker network
atlasmart-cassandra, datacenter dc1,
racks rack1..rack3, and 16 virtual nodes per
node. Chapter 09 uses keyspace
atlasmart_writepath with
NetworkTopologyStrategy, replication factor (RF)
3, and LOCAL_QUORUM unless an exercise
deliberately changes consistency level. New tables explicitly
use UnifiedCompactionStrategy (UCS); table TTL defaults to
zero and gc_grace_seconds is not changed.
Authentication, client TLS, internode TLS, and remote JMX are
disabled only inside the isolated learning network. Apache
Cassandra Java Driver 4.19.3 is optional for routing examples;
every mandatory exercise remains free/local with
cqlsh, nodetool, Docker, and shell
commands.
Run commands only against the disposable Apache Cassandra course lab or another explicitly approved non-production environment. Confirm node, keyspace, table, container, volume, path, and datacenter targets before destructive, failure-injection, cleanup, repair, restore, security, or topology operations. Capture current state and expected rollback/recovery evidence first; output and timings can differ by host, operating system, Java runtime, Docker/runtime, driver, and Cassandra configuration.
1. A coordinator routes work; it is not a permanent primary
A Cassandra client connects to one node for a request. That node is the coordinator for that request. A modern token-aware driver tries to choose a coordinator intelligently—often preferring a replica in the local datacenter—but the coordinator role is temporary and request-scoped. The actual data placement is determined from the partition key, its token, the keyspace replication strategy, and topology metadata.
For order_id='order-1001', every clustering row in
that partition maps to the same token. With RF=3 in this
three-node lab, all three nodes are natural replicas for the
partition, but that does not make any one of them a permanent
leader. The coordinator forwards the mutation to the replicas
and waits for enough replica acknowledgments to satisfy the
requested consistency level (CL).
| Term | What it controls | What it does not mean |
|---|---|---|
| Coordinator | Routes one request, collects replica responses, applies the requested CL. | It is not the only durable owner and is not a permanent primary. |
| Replica | Stores a copy of the partition according to placement. | Every replica need not acknowledge before every CL can succeed. |
| RF | Number of replicas Cassandra intends to maintain per datacenter. | RF is not a read/write CL and is not a backup count. |
| CL | How many relevant replica responses are required for this operation to be considered successful. | Success does not prove all replicas are already equal. |
| Token-aware routing | Lets a driver prefer nodes that own the partition and respect local-DC policy. | It does not change the replica set or bypass consistency rules. |
“The coordinator receives the write, so the coordinator stores the only authoritative copy” is false. A coordinator can even be a non-replica. Production troubleshooting must distinguish coordinator latency from replica storage latency and network delay.
2. Consistency is acknowledgment math, not a vague durability adjective
In this one-datacenter RF=3 lab,
LOCAL_QUORUM requires two acknowledgments from
replicas in dc1. ONE requires one.
ALL requires all three. QUORUM is also
two here, but in a multi-datacenter keyspace it is computed
across the total replication factor, whereas
LOCAL_QUORUM deliberately scopes its requirement to
the local datacenter.
| CL for RF=3, one DC | Required acknowledgments | Failure implication |
|---|---|---|
| ONE | 1 | Can succeed with two replicas unavailable; weakest immediate replica coverage. |
| LOCAL_QUORUM | 2 local replicas | Survives one replica loss in dc1 while preserving a quorum intersection with LOCAL_QUORUM reads. |
| QUORUM | 2 replicas | Same number in this topology, but semantics differ once multiple DCs exist. |
| ALL | 3 replicas | Cannot succeed if any replica is unavailable or fails to acknowledge in time. |
A successful CL response means the coordinator received enough acknowledgments for that operation. It does not prove that every intended replica received the mutation, that a backup exists, that repair is unnecessary, or that the commit log on each acknowledging replica was physically fsynced under every commit-log sync mode. Those are separate durability and convergence questions.
3. Lab: identify replicas and trace one write
# Verify the shared lab if it already exists.docker exec atlasmart-cass-1 nodetool versiondocker exec atlasmart-cass-1 nodetool status# Standalone recreation path. Skip resources that already exist.docker network inspect atlasmart-cassandra >/dev/null 2>&1 || docker network create atlasmart-cassandradocker volume create atlasmart-cass-1-datadocker volume create atlasmart-cass-2-datadocker volume create atlasmart-cass-3-datadocker inspect atlasmart-cass-1 >/dev/null 2>&1 || docker run -d --name atlasmart-cass-1 --hostname atlasmart-cass-1 --network atlasmart-cassandra -e CASSANDRA_CLUSTER_NAME=atlasmart-course -e CASSANDRA_DC=dc1 -e CASSANDRA_RACK=rack1 -e CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch -e CASSANDRA_NUM_TOKENS=16 -v atlasmart-cass-1-data:/var/lib/cassandra cassandra:5.0.9# Wait until node 1 answers before starting peers.docker exec atlasmart-cass-1 nodetool statusdocker inspect atlasmart-cass-2 >/dev/null 2>&1 || docker run -d --name atlasmart-cass-2 --hostname atlasmart-cass-2 --network atlasmart-cassandra -e CASSANDRA_CLUSTER_NAME=atlasmart-course -e CASSANDRA_DC=dc1 -e CASSANDRA_RACK=rack2 -e CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch -e CASSANDRA_NUM_TOKENS=16 -e CASSANDRA_SEEDS=atlasmart-cass-1 -v atlasmart-cass-2-data:/var/lib/cassandra cassandra:5.0.9docker inspect atlasmart-cass-3 >/dev/null 2>&1 || docker run -d --name atlasmart-cass-3 --hostname atlasmart-cass-3 --network atlasmart-cassandra -e CASSANDRA_CLUSTER_NAME=atlasmart-course -e CASSANDRA_DC=dc1 -e CASSANDRA_RACK=rack3 -e CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch -e CASSANDRA_NUM_TOKENS=16 -e CASSANDRA_SEEDS=atlasmart-cass-1 -v atlasmart-cass-3-data:/var/lib/cassandra cassandra:5.0.9# Wait until all three nodes show UN in dc1 before continuing.docker exec atlasmart-cass-1 nodetool status# Open cqlsh on node 1 for the SQL/CQL blocks that follow.docker exec -it atlasmart-cass-1 cqlsh
CREATE KEYSPACE IF NOT EXISTS atlasmart_writepathWITH replication = {'class':'NetworkTopologyStrategy','dc1':3};CREATE TABLE IF NOT EXISTS atlasmart_writepath.order_write_events ( order_id text, event_time timeuuid, status text, source text, note text, expires_at timestamp, PRIMARY KEY ((order_id), event_time)) WITH CLUSTERING ORDER BY (event_time DESC) AND compaction = {'class':'UnifiedCompactionStrategy'};CONSISTENCY LOCAL_QUORUM;SELECT keyspace_name, replicationFROM system_schema.keyspacesWHERE keyspace_name = 'atlasmart_writepath';
# getendpoints asks Cassandra which endpoints own the partition key.docker exec atlasmart-cass-1 nodetool getendpoints atlasmart_writepath order_write_events order-1001# Confirm rack/DC placement and all nodes are currently Up/Normal.docker exec atlasmart-cass-1 nodetool status atlasmart_writepath
CONSISTENCY LOCAL_QUORUM;TRACING ON;INSERT INTO atlasmart_writepath.order_write_events(order_id, event_time, status, source, note)VALUES ('order-1001', now(), 'PAID', 'checkout-api', 'first traced write');TRACING OFF;SELECT order_id, event_time, status, source, WRITETIME(status)FROM atlasmart_writepath.order_write_eventsWHERE order_id='order-1001';
The trace is evidence about the coordinator path for that
particular request. Exact event lines, node addresses, token
values, and timings depend on the runtime. Record them rather
than copying a canned trace.
nodetool getendpoints is placement evidence; the
trace is request evidence; neither alone proves on-disk
persistence on every replica.
4. Controlled replica loss: quorum success versus ALL failure
Blast radius: only atlasmart-cass-3 in the
disposable Docker cluster is stopped. Do not perform this on
unrelated local data or production. Wait until the remaining
nodes report it down so the outcome is deterministic;
immediately after a failure, the request may instead wait and
time out while failure detection converges.
docker stop atlasmart-cass-3# Re-run until atlasmart-cass-3 is shown Down/Normal (DN) by the live nodes.docker exec atlasmart-cass-1 nodetool status atlasmart_writepath# After the CQL exercises below:docker start atlasmart-cass-3# Wait until all three nodes are Up/Normal (UN) again.docker exec atlasmart-cass-1 nodetool status atlasmart_writepath
CONSISTENCY LOCAL_QUORUM;INSERT INTO atlasmart_writepath.order_write_events(order_id, event_time, status, source, note)VALUES ('order-1001', now(), 'PACKING', 'warehouse-api', 'quorum can still succeed');CONSISTENCY ALL;INSERT INTO atlasmart_writepath.order_write_events(order_id, event_time, status, source, note)VALUES ('order-1001', now(), 'SHOULD_FAIL', 'lab', 'ALL needs all three replicas');CONSISTENCY LOCAL_QUORUM;
Once node 3 is down, the LOCAL_QUORUM write can succeed with two acknowledgments; the ALL write cannot satisfy three. Cassandra may preserve a hint for the missed replica, and repair remains the authoritative anti-entropy mechanism. The expected error class can be unavailable after the node is known down, or a write timeout if the failure is detected after the request begins. Both are useful evidence; neither should be hidden behind a generic “write failed” message.
docker exec atlasmart-cass-1 nodetool listpendinghintsdocker exec atlasmart-cass-2 nodetool listpendinghints# Output is topology/timing dependent; a zero count does not disprove the write path.# Hints may already have been delivered after node 3 returns.
5. Production judgment
Token-aware local routing reduces avoidable network hops but does not change the consistency contract. Choose RF and CL from explicit failure tolerance, read/write intersection, latency, and recovery requirements. Watch coordinator p95/p99 write latency, replica local write latency, timeouts/unavailable errors, dropped mutations, hint pressure, pending flushes, disk saturation, and network health together. A retry policy must understand idempotency: a timeout can mean “unknown whether enough replicas applied the mutation,” not “nothing happened.”
Do not tune by forcing the coordinator to a favored node, by assuming CL ONE is “fast enough,” or by treating RF=3 as three backups. The next lesson moves inside each replica and shows what its acknowledgment is actually protecting through the commit log and memtable.
Verification checklist
- All three nodes return to UN after the failure exercise.
- The keyspace reports RF=3 in dc1.
-
nodetool getendpointsreturns the intended replica set for the partition. - A LOCAL_QUORUM write succeeds with one replica down, while ALL cannot meet its requirement.
- You saved your own trace/error output instead of treating sample output as measured evidence.
Check your understanding
- Can a non-replica coordinate a write?
- For RF=3 in one DC, how many acknowledgments does LOCAL_QUORUM require?
- Does LOCAL_QUORUM success prove all three replicas contain the mutation immediately?
- Why can a stopped replica produce either an unavailable error or a write timeout?
- What does token-aware routing optimize?
Review the answers
1. Yes. The coordinator is the contacted/request-routing node; placement decides which nodes are replicas.
2. Two local replica acknowledgments.
3. No. It proves the required quorum acknowledged; another replica may be unavailable and catch up through hints or repair.
4. The result depends on whether failure detection already classified the replica as unavailable before the request versus the request waiting for acknowledgments that never arrive.
5. Coordinator choice and locality; it does not change token ownership, RF, or CL.
Summary and next bridge
A write begins as a routed request, becomes coordinator work, fans out to natural replicas, and succeeds only when the chosen CL receives enough replica acknowledgments. Next we step inside one replica to separate commit-log durability from memtable state and understand exactly what can be acknowledged before an SSTable exists.
Authoritative references
- Apache Cassandra 5.0 — Storage Engine / write path, commit log, memtables, flushes
- Apache Cassandra 5.0 — Hinted handoff
- Apache Cassandra CQL — DML, TIMESTAMP, TTL, WRITETIME
- nodetool tablestats
- nodetool listpendinghints
- Apache Cassandra native protocol — write timeout/failure metadata
- Apache Cassandra downloads and current releases
- Apache Cassandra Java Driver