Chapter 10 · The Read Path: Bloom Filters, Indexes, SSTables, Caches, Merging, and Reconciliation
Coordinator Read Requests, Replica Selection, Data / Digest Responses, and Consistency
Follow an AtlasMart read through coordinator selection, replica data/digest responses, consistency-level completion, mismatch detection, and reconciliation.
Learning outcomes
AtlasMart's order page sometimes returns quickly and sometimes waits. Before blaming disk, the team must know which node coordinated the read, which replicas participated, which consistency level was requested, and whether replica results agreed. This lesson traces that distributed half of the read path first.
Separate driver routing, request coordination, replica ownership, and consistency-level completion.
Explain full data versus digest responses and what a digest mismatch changes.
Derive response requirements for LOCAL_ONE, LOCAL_QUORUM, QUORUM, and ALL at RF=3 in one DC.
Use tracing and nodetool topology evidence during a reversible one-replica failure.
State what read success proves and does not prove about every replica.
The mandatory labs continue the established AtlasMart
disposable cluster: cassandra:5.0.9, cluster
atlasmart-course, Docker network
atlasmart-cassandra, nodes
atlasmart-cass-1..3, datacenter dc1,
racks rack1..rack3, and 16 virtual nodes per
node. Chapter 10 uses atlasmart_readpath with
NetworkTopologyStrategy, replication factor (RF)
3, and normally LOCAL_QUORUM. 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 remain disabled only inside the isolated learning network.
The Apache Cassandra Java Driver 4.19.3 is optional for
client-routing discussion; all mandatory evidence uses free
local cqlsh/nodetool. Capture
nodetool version, cqlsh --version,
and java -version on your machine rather than
treating a prose baseline as runtime proof.
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.
Core terms before tracing a read
A coordinator is the Cassandra node handling one client request; it is not a permanent leader. A replica stores a copy of a partition according to replication placement. A partition is the rows sharing a partition key; the partitioner hashes that key to a token, which participates in replica ownership. A consistency level (CL) specifies how many appropriately scoped replica responses a coordinator needs for the operation. An SSTable (Sorted String Table) is immutable on-disk storage. A tombstone is a distributed deletion/expiration marker. A Bloom filter is a probabilistic membership structure that can say “definitely absent” or “possibly present”; for its intended membership check it can produce false positives but not false negatives. Reconciliation merges cell versions and tombstones to construct the newest visible result. Read repair may write that reconciled result back to stale replicas involved in a request; it does not replace scheduled anti-entropy repair.
1. Driver → coordinator → natural replicas
A token-aware driver can choose a coordinator near the natural
replicas, but the contacted node is still only the coordinator
for that request. It identifies the replica set and sends read
commands. With RF=3 and LOCAL_QUORUM in one
datacenter, two local replica responses are sufficient to
satisfy the client-visible CL.
Cassandra can request a full data response from one replica and
digest responses—hashes of the selected result—from other
replicas needed by the request. If digests match, duplicate full
rows need not cross the network. If they differ, the coordinator
must obtain enough full data to reconcile the newest visible
cells/tombstones. The table's read_repair setting
then governs whether stale replicas involved in the read are
updated as part of the foreground request.
| Stage | Evidence | Proves | Does not prove |
|---|---|---|---|
| Coordinator | trace / client endpoint | which node handled request | that coordinator stores the only copy |
| Replica set | topology/getendpoints | natural replica ownership | which will answer fastest |
| Data/digest | trace | which endpoints/messages participated | that every RF replica returned full data |
| CL success | LOCAL_QUORUM result | enough scoped responses arrived | all replicas are identical |
| Mismatch | trace/reconciliation | involved results differ | scheduled repair is unnecessary |
2. Reproduce and trace
docker exec atlasmart-cass-1 nodetool versiondocker exec atlasmart-cass-1 nodetool statusdocker exec atlasmart-cass-1 java -version# Recreate only if the shared course lab does not 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 is UN 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.9docker exec atlasmart-cass-1 nodetool statusdocker exec -it atlasmart-cass-1 cqlsh
CREATE KEYSPACE IF NOT EXISTS atlasmart_readpathWITH replication = {'class':'NetworkTopologyStrategy','dc1':3};CREATE TABLE IF NOT EXISTS atlasmart_readpath.order_reads_by_customer ( customer_id text, order_month date, order_time timestamp, order_id uuid, status text, total decimal, note text, PRIMARY KEY ((customer_id,order_month),order_time,order_id)) WITH CLUSTERING ORDER BY (order_time DESC,order_id ASC) AND compaction = {'class':'UnifiedCompactionStrategy'} AND caching = {'keys':'ALL','rows_per_partition':'NONE'};CONSISTENCY LOCAL_QUORUM;INSERT INTO atlasmart_readpath.order_reads_by_customer(customer_id,order_month,order_time,order_id,status,total,note)VALUES ('cust-42','2026-09-01','2026-09-07T16:00:00Z',00000000-0000-0000-0000-000000000042,'PAID',129.90,'baseline');SELECT * FROM atlasmart_readpath.order_reads_by_customerWHERE customer_id='cust-42' AND order_month='2026-09-01';TRACING ON;SELECT status,total,note FROM atlasmart_readpath.order_reads_by_customerWHERE customer_id='cust-42' AND order_month='2026-09-01';TRACING OFF;
CREATE TABLE IF NOT EXISTS atlasmart_readpath.read_probe_by_key ( probe_key text PRIMARY KEY, value text) WITH compaction = {'class':'UnifiedCompactionStrategy'};INSERT INTO atlasmart_readpath.read_probe_by_key (probe_key,value) VALUES ('probe-42','visible');
docker exec atlasmart-cass-1 nodetool statusdocker exec atlasmart-cass-1 nodetool getendpoints atlasmart_readpath read_probe_by_key probe-42
Exact trace lines, endpoints, and microseconds depend on token placement, snitch latency history, speculative retry, and runtime state. Capture them; do not fabricate fixed output.
3. One replica unavailable: RF and CL are different contracts
Pause only node 3 in this disposable lab. With RF=3,
LOCAL_QUORUM can still succeed because two local
replicas remain. ALL cannot. An
unavailable error means Cassandra already knows there
are insufficient live replicas; a timeout means enough
may be alive, but sufficient responses did not arrive before the
deadline.
docker pause atlasmart-cass-3docker exec atlasmart-cass-1 nodetool status
CONSISTENCY LOCAL_QUORUM;SELECT * FROM atlasmart_readpath.read_probe_by_key WHERE probe_key='probe-42';CONSISTENCY ALL;SELECT * FROM atlasmart_readpath.read_probe_by_key WHERE probe_key='probe-42';CONSISTENCY LOCAL_ONE;SELECT * FROM atlasmart_readpath.read_probe_by_key WHERE probe_key='probe-42';
docker unpause atlasmart-cass-3docker exec atlasmart-cass-1 nodetool status
CL specifies response requirements, not a rule that every replica returns full data. Replica selection, digest/data requests, speculation, mismatch handling, and topology determine actual messages. Use trace evidence.
4. Verification and concept checks
- All three nodes return
UNafter reset. -
RF=3 and
LOCAL_QUORUMare recorded explicitly. - The trace is interpreted without hard-coding timing.
- Unavailable, timeout, digest mismatch, reconciliation, and repair are distinguished.
Check your understanding
- Why can a read return without full data from every replica?
- At RF=3 in one DC, how many responses does LOCAL_QUORUM require?
- Does LOCAL_QUORUM success prove the third replica is identical?
- What does a digest mismatch mean?
- Unavailable versus timeout?
Review the answers
1. Digest responses can confirm agreement while a full data response supplies the result; exact messaging is version/runtime dependent.
2. Two local replica responses.
3. No. It proves the requested CL was satisfied.
4. Replica results for the read differ, so the coordinator needs full data sufficient for reconciliation.
5. Unavailable is known insufficient live replicas; timeout is insufficient responses before the deadline despite the possibility enough replicas are alive.
Production judgment
Read latency is not a single disk or cache number. Evaluate partition rows/bytes and skew, clustering slices, RF and CL, replica locality/health, p95/p99 latency, SSTables touched per logical read, tombstones scanned, compaction state, Bloom false positives, key/chunk/page-cache state, disk throughput/queueing, JVM/GC, driver timeout/retry/speculation, repair state, and failure-domain health. Large partitions or overlapping SSTables can dominate an otherwise healthy cache. Conversely, aggressive cache allocation can steal memory from the OS page cache or other off-heap structures.
Do not generalize a warm laptop result. Record the Cassandra patch, SSTable format, topology, RF/CL, dataset and partition distribution, payload/result size, TTL/delete rate, compaction state, concurrent load, disk/network, warmup, and failure injection. Forced flush/compaction is a controlled lab/maintenance action, not routine performance tuning. Lesson 2 moves inside one replica and follows Bloom/index/cache checks to the SSTable bytes.
Summary and next bridge
A Cassandra read is distributed coordination before it is storage I/O. The coordinator gathers enough replica evidence for the CL, compares results, and reconciles differences. Next, follow the same key through Bloom filters, indexes, caches, and SSTable storage.
Authoritative references
Re-check these version-sensitive sources when regenerating the lesson.