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.

Intermediate100–140 minutesCoordinator + replica routing labApache Cassandra 5.0.9 · cqlsh/nodetool · Java Driver 4.19.3 optional · RF=3 · UCSLast reviewed: September 2026

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.

01

Explain the roles of client driver, coordinator, natural replica, token, and consistency level without treating them as synonyms.

02

Relate a partition key to token ownership and the replica set selected by NetworkTopologyStrategy.

03

Calculate the acknowledgments required by ONE, LOCAL_QUORUM, QUORUM, and ALL for RF=3 in one datacenter.

04

Use cqlsh tracing and nodetool evidence to distinguish coordinator work from replica work.

05

Run a reversible one-replica failure experiment and state what a successful quorum write does and does not guarantee.

Chapter 09 lab baseline

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.

Execution and safety note

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.
Wrong mental model

“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

bash · verify or recreate the disposable three-node cluster
# 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
sql · create the isolated write-path keyspace and table
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';
bash · inspect replica ownership for one partition key
# 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
sql · trace a LOCAL_QUORUM write from cqlsh
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.

bash · stop one replica, observe state, then restore it
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
sql · compare LOCAL_QUORUM and ALL while one replica is down
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.

bash · inspect pending hints after the controlled failure
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 getendpoints returns 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

  1. Can a non-replica coordinate a write?
  2. For RF=3 in one DC, how many acknowledgments does LOCAL_QUORUM require?
  3. Does LOCAL_QUORUM success prove all three replicas contain the mutation immediately?
  4. Why can a stopped replica produce either an unavailable error or a write timeout?
  5. 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

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.