Chapter 09 · The Write Path: Coordinator, Commit Log, Memtables, Replicas, and Acknowledgment
Trace a Write Across Replicas and Relate Acknowledgments to Consistency Level
Run a controlled replica-loss game day and connect ONE/LOCAL_QUORUM/ALL client outcomes to replica acknowledgments, hints, recovery, and repair.
Learning outcomes
AtlasMart is preparing a production readiness review. The team can name every write-path component, but the final requirement is operational: demonstrate what a client sees at different consistency levels while one replica disappears, correlate the result with tracing/hints/table statistics, recover the node, and explain which guarantees still require repair and monitoring.
Run a repeatable one-replica failure drill at ONE, LOCAL_QUORUM, and ALL without unsafe host-network manipulation.
Relate successful responses and native-protocol errors to received versus required replica acknowledgments.
Inspect hints and node/table statistics while distinguishing best-effort catch-up from full anti-entropy repair.
Explain timeout ambiguity and why retries require idempotency discipline.
Produce a write-path evidence checklist suitable for a production runbook or design review.
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. The end-to-end write-path contract
The application issues a mutation through the native protocol. A driver selects a node using topology/routing metadata; that request-scoped coordinator identifies natural replicas for the partition and sends mutations. Each responding replica performs its local commit-log/memtable work according to configuration. The coordinator returns success when the requested CL has enough acknowledgments—or returns an unavailable, timeout, or failure response when it cannot meet that contract.
| Layer | Evidence to capture | Failure question |
|---|---|---|
| Client/driver | Requested CL, timeout, retry/idempotency, coordinator if exposed | Was the outcome definitively failed, or unknown after timeout? |
| Coordinator | Trace events, received/block-for metadata on errors, proxy latency | Did coordination/network wait dominate? |
| Replica | UN/DN state, local write latency, memtable/pending flushes | Could the replica accept and persist the mutation fast enough? |
| Catch-up | Pending hints, repair history/runbook | How will a missed replica converge after recovery? |
| Storage | Commit-log mode, SSTable/flush state, disk/GC metrics | What local durability and capacity assumptions support each ack? |
2. Lab preparation: clean evidence and known replica health
# 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';
docker exec atlasmart-cass-1 nodetool status atlasmart_writepathdocker exec atlasmart-cass-1 nodetool getendpoints atlasmart_writepath order_write_events order-drilldocker exec atlasmart-cass-1 nodetool tablestats atlasmart_writepath.order_write_eventsdocker exec atlasmart-cass-1 nodetool listpendinghintsdocker exec atlasmart-cass-2 nodetool listpendinghints
CONSISTENCY ALL;TRACING ON;INSERT INTO atlasmart_writepath.order_write_events(order_id,event_time,status,source,note)VALUES ('order-drill',now(),'BASELINE','game-day','all replicas available');TRACING OFF;CONSISTENCY LOCAL_QUORUM;
Save the baseline trace. Exact node addresses and timings are your evidence. If ALL does not succeed while all nodes are healthy, stop the drill and fix the lab before injecting failure.
3. Stop one replica and compare consistency levels
Use Docker stop, not host firewall rules or clock skew. Wait for the surviving nodes to classify node 3 as down. This keeps the blast radius contained and the reset path obvious.
docker stop atlasmart-cass-3# Repeat until node 3 is shown DN by a live node.docker exec atlasmart-cass-1 nodetool status atlasmart_writepath
CONSISTENCY ONE;TRACING ON;INSERT INTO atlasmart_writepath.order_write_events(order_id,event_time,status,source,note)VALUES ('order-drill',now(),'ONE_OK','game-day','one ack required');TRACING OFF;CONSISTENCY LOCAL_QUORUM;TRACING ON;INSERT INTO atlasmart_writepath.order_write_events(order_id,event_time,status,source,note)VALUES ('order-drill',now(),'QUORUM_OK','game-day','two local acks required');TRACING OFF;CONSISTENCY ALL;TRACING ON;INSERT INTO atlasmart_writepath.order_write_events(order_id,event_time,status,source,note)VALUES ('order-drill',now(),'ALL_EXPECTED_TO_FAIL','game-day','three acks required');TRACING OFF;CONSISTENCY LOCAL_QUORUM;
With RF=3 and one replica known down, ONE and LOCAL_QUORUM have enough remaining replicas; ALL does not. Immediately after an ungraceful failure, a request can instead time out while failure detection catches up. The native protocol includes write-timeout fields such as the CL, number of acknowledgments received, number required (block-for), and write type. Preserve those details in application logs rather than flattening every exception into “Cassandra unavailable.”
4. Inspect catch-up evidence, then recover the node
docker exec atlasmart-cass-1 nodetool listpendinghintsdocker exec atlasmart-cass-2 nodetool listpendinghintsdocker start atlasmart-cass-3# Wait until all nodes show UN.docker exec atlasmart-cass-1 nodetool status atlasmart_writepath# Re-check hints; they may already have drained by the time you inspect.docker exec atlasmart-cass-1 nodetool listpendinghintsdocker exec atlasmart-cass-2 nodetool listpendinghintsdocker exec atlasmart-cass-1 nodetool tablestats atlasmart_writepath.order_write_eventsdocker exec atlasmart-cass-3 nodetool tablestats atlasmart_writepath.order_write_events
CONSISTENCY ALL;SELECT order_id,event_time,status,source,noteFROM atlasmart_writepath.order_write_eventsWHERE order_id='order-drill';CONSISTENCY LOCAL_QUORUM;
A successful ALL read after recovery shows all replicas responded sufficiently for that read and the reconciled result is visible; it is not a substitute for the cluster’s scheduled repair strategy. Hints are best-effort and time-bounded. Production runbooks still need repair to synchronize missed ranges and verify long-lived replica convergence.
5. The retry trap: timeout means unknown outcome
A write timeout is not equivalent to a transaction rollback. Some replicas may already have applied the mutation even though the coordinator did not collect enough acknowledgments before the deadline. Retrying a naturally idempotent upsert with the same intended value may be safe; retrying counters, non-idempotent read-before-write workflows, or business side effects can duplicate work. Driver retry policies therefore belong in the data model and API correctness review, not only the networking configuration.
A larger timeout changes how long the client waits; it does not increase disk throughput, repair a failed replica, or guarantee more acknowledgments. First identify whether the cause is topology/failure detection, commit-log/storage latency, flush pressure, network delay, GC, oversized mutations, or overload.
6. Production acceptance checklist
| Area | Evidence required before production |
|---|---|
| Routing | Local-DC/token-aware driver policy tested; contact points are discovery seeds, not a fixed coordinator list. |
| Consistency | RF/CL matrix tied to failure modes and read-after-write requirements. |
| Durability | commitlog_sync mode, storage device/fsync assumptions, RF/CL, and crash model documented. |
| Capacity | p50/p95/p99/max write latency, pending flushes, disk/commit-log headroom, GC, network, mutation size and concurrency measured under representative load. |
| Convergence | Hints monitored, repair schedule tested, node replacement/recovery runbooks exercised. |
| Client correctness | Timeout/retry/idempotency policy tested with ambiguous outcomes and duplicate delivery. |
| Security | Authentication, authorization, TLS, JMX/network isolation, secrets and audit controls added beyond this isolated lab. |
Verification checklist
- All three nodes end UN.
- ONE and LOCAL_QUORUM succeed with one replica down; ALL cannot meet RF=3.
- You captured at least one trace and the exact failure/timeout class from your environment.
- You inspected pending hints before/after recovery without assuming a particular count.
- An ALL read succeeds after recovery.
- You can state why repair remains necessary even if hints appear to drain.
Check your understanding
- Why can LOCAL_QUORUM succeed with one replica down at RF=3 in one DC?
- What does a write timeout tell you about whether the mutation applied?
- Are hints a replacement for repair?
- Why should applications preserve received/block-for details?
- What is the bridge to Chapter 10?
Review the answers
1. It requires two local replica acknowledgments, and two replicas remain available.
2. The outcome can be ambiguous: some replicas may have applied it even though the coordinator did not receive enough acknowledgments in time.
3. No. Hints are best-effort catch-up; repair is the anti-entropy mechanism that compares replica data and streams differences.
4. They distinguish how many acknowledgments arrived versus how many the requested CL required, which improves diagnosis of topology, latency, and failure behavior.
5. Writes create versions in memtables/SSTables; reads must locate those versions, merge them, reconcile replicas, and account for tombstones, caches, Bloom filters and SSTable amplification.
Summary and next bridge
Chapter 09 traced the complete write path: routing selects a coordinator, placement selects replicas, each replica records commit-log/memtable state, flush creates SSTables, timestamps and tombstones define reconciliation, and CL defines the client acknowledgment threshold. Chapter 10 turns the direction around: the coordinator must now find and reconcile the newest visible data across memtables, indexes, Bloom filters, SSTables, caches, replicas, and tombstones.
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