A Cassandra read looks simple from the client: one CQL statement, one row back. Underneath, three separate machines do work you can tune independently. A coordinator decides which replicas to ask and how long to wait. Each replica walks a storage engine that may have to consult a dozen files to assemble one partition. Then the coordinator compares what the replicas said and, if they disagree, repairs them before answering. A p99 regression almost always lives in exactly one of those stages, and the fix for one stage does nothing for the others.

This article follows a single LOCAL_QUORUM read at replication factor three from the driver to the disk and back, using the behaviour of Cassandra 4.x and 5.0. It corrects a few things commonly repeated about the read path (the coordinator does not ask every replica for full data, and read repair at quorum is not a background job), then shows how to trace a slow query, which metrics separate the three stages, and what to change when each stage is the bottleneck.

One LOCAL_QUORUM read, RF=3: plan, fetch, reconcileClient drivertoken-awareCoordinatortoken, replicas, blockFor=2Replica A: DATA requestReplica B: DIGEST requestReplica C: only if speculative retrydigests match: return Amismatch: full data from all,merge, blocking read repairInside each replica (single-partition read)Row cacheoff by defaultMemtableslive + flushingSSTable filtertimestamps, boundsBloom filterper SSTablePartition indexBIG or BTIData filechunk cache, diskMerge iteratornewest timestamp wins per cell; tombstones shadow older data; then filter, limit, pageLatency = slowest replica you wait for + work per SSTable touched + tombstones scanned.Every tuning knob in this article moves one of those three terms.
The coordinator asks one replica for data and the rest for digests; each replica narrows its SSTables before merging every source by timestamp.

The coordinator plans the read

The driver usually picks the coordinator. A token-aware load balancing policy hashes the partition key with the same Murmur3 partitioner the cluster uses and sends the request straight to a replica, which saves a network hop and means the coordinator can serve its own copy locally. Without token awareness, any node in the local datacenter coordinates and forwards everything.

The coordinator computes the token, looks up the replicas for that token in the keyspace's replication strategy, and works out blockFor, the number of responses the consistency level requires: two for LOCAL_QUORUM with RF=3 in the local datacenter. It then orders the replicas by the dynamic snitch, which scores nodes on recent latency, and sends exactly one full data request to the best replica and digest requests to the others it needs. A digest is a hash of the result the replica would have returned. That is the main bandwidth saving on the read path: for a 50 KB partition at quorum you move one 50 KB payload plus a few bytes, not two.

If the chosen replicas are slow, speculative retry kicks in. The table option speculative_retry defaults to 99p: when a replica has not answered within that table's observed 99th percentile read latency, the coordinator sends one extra request to a replica it has not used yet. It also accepts fixed values such as '50ms', 'ALWAYS' and 'NONE'. Speculation trims the tail caused by one sick node at the cost of a little extra load; it does nothing when all replicas are slow.

# Coordinator read, simplified (single partition, no paging)
def coordinate_read(query, cl):
    token    = murmur3(query.partition_key)
    replicas = snitch.sort_by_proximity(replicas_for(token, cl.local_dc_only))
    need     = block_for(cl, replicas)              # LOCAL_QUORUM, RF=3 -> 2
    targets  = replicas[:need]
    futures  = [send_data(targets[0], query)] + [send_digest(r, query) for r in targets[1:]]
    if not wait_all(futures, timeout=table.speculation_threshold):    # e.g. table p99
        spare = first_unused(replicas, targets)
        if spare: futures.append(send_data_or_digest(spare, query))
    responses = wait_for(need, futures, timeout=read_request_timeout)  # else ReadTimeoutException
    if all_digests_match(responses):
        return responses.data()
    full = [send_data(r, query) for r in responded(responses)]       # second round trip
    merged = reconcile(wait_all(full))                               # newest timestamp wins
    blocking_read_repair(merged, full)                               # writes back before replying
    return merged

Inside one replica: from memtable to disk

On each replica the read becomes a single-partition read command. If the row cache is enabled for the table (it is off by default and rarely worth it for anything but small, hot, read-mostly partitions) a hit returns immediately. Otherwise the replica builds a merge over every source that might hold part of the partition: the live memtable, any memtables currently flushing, and the SSTables on disk.

SSTable selection is where most reads are won or lost. The replica first discards SSTables using metadata it keeps in memory: minimum and maximum clustering values tell it a file cannot contain the requested slice, and for queries that name specific cells the engine visits files in descending order of maximum timestamp and can stop once it has a value newer than anything an older file could hold. Each surviving SSTable's Bloom filter then answers whether the partition key might be present. A negative is certain; a positive is wrong with the probability set by bloom_filter_fp_chance (0.01 by default, 0.1 for leveled compaction).

Files that pass the filter must locate the partition. In the long-standing BIG format that means an in-memory partition summary sampling the on-disk partition index, the key cache short-circuiting both, and a row index inside large partitions so the engine can seek to a clustering range without reading the whole partition. Cassandra 5.0 adds the BTI format, which replaces the summary and index with on-disk tries that are cheap to search without a key cache. Either way the result is a file offset. The engine reads the compressed chunk containing it through the chunk cache or the OS page cache, decompresses the whole chunk, and hands rows to the merge. A large compression chunk length makes every small read pay to decompress the full chunk, which is why read-heavy tables with small rows often benefit from a smaller chunk_length_in_kb.

The coordinator plans the read

The driver usually picks the coordinator. A token-aware load balancing policy hashes the partition key with the same Murmur3 partitioner the cluster uses and sends the request straight to a replica, which saves a network hop and means the coordinator can serve its own copy locally. Without token awareness, any node in the local datacenter coordinates and forwards everything.

The coordinator computes the token, looks up the replicas for that token in the keyspace's replication strategy, and works out blockFor, the number of responses the consistency level requires: two for LOCAL_QUORUM with RF=3 in the local datacenter. It then orders the replicas by the dynamic snitch, which scores nodes on recent latency, and sends exactly one full data request to the best replica and digest requests to the others it needs. A digest is a hash of the result the replica would have returned. That is the main bandwidth saving on the read path: for a 50 KB partition at quorum you move one 50 KB payload plus a few bytes, not two.

If the chosen replicas are slow, speculative retry kicks in. The table option speculative_retry defaults to 99p: when a replica has not answered within that table's observed 99th percentile read latency, the coordinator sends one extra request to a replica it has not used yet. It also accepts fixed values such as '50ms', 'ALWAYS' and 'NONE'. Speculation trims the tail caused by one sick node at the cost of a little extra load; it does nothing when all replicas are slow.

# Coordinator read, simplified (single partition, no paging)
def coordinate_read(query, cl):
    token    = murmur3(query.partition_key)
    replicas = snitch.sort_by_proximity(replicas_for(token, cl.local_dc_only))
    need     = block_for(cl, replicas)              # LOCAL_QUORUM, RF=3 -> 2
    targets  = replicas[:need]
    futures  = [send_data(targets[0], query)] + [send_digest(r, query) for r in targets[1:]]
    if not wait_all(futures, timeout=table.speculation_threshold):    # e.g. table p99
        spare = first_unused(replicas, targets)
        if spare: futures.append(send_data_or_digest(spare, query))
    responses = wait_for(need, futures, timeout=read_request_timeout)  # else ReadTimeoutException
    if all_digests_match(responses):
        return responses.data()
    full = [send_data(r, query) for r in responded(responses)]       # second round trip
    merged = reconcile(wait_all(full))                               # newest timestamp wins
    blocking_read_repair(merged, full)                               # writes back before replying
    return merged

Inside one replica: from memtable to disk

On each replica the read becomes a single-partition read command. If the row cache is enabled for the table (it is off by default and rarely worth it for anything but small, hot, read-mostly partitions) a hit returns immediately. Otherwise the replica builds a merge over every source that might hold part of the partition: the live memtable, any memtables currently flushing, and the SSTables on disk.

SSTable selection is where most reads are won or lost. The replica first discards SSTables using metadata it keeps in memory: minimum and maximum clustering values tell it a file cannot contain the requested slice, and for queries that name specific cells the engine visits files in descending order of maximum timestamp and can stop once it has a value newer than anything an older file could hold. Each surviving SSTable's Bloom filter then answers whether the partition key might be present. A negative is certain; a positive is wrong with the probability set by bloom_filter_fp_chance (0.01 by default, 0.1 for leveled compaction).

Files that pass the filter must locate the partition. In the long-standing BIG format that means an in-memory partition summary sampling the on-disk partition index, the key cache short-circuiting both, and a row index inside large partitions so the engine can seek to a clustering range without reading the whole partition. Cassandra 5.0 adds the BTI format, which replaces the summary and index with on-disk tries that are cheap to search without a key cache. Either way the result is a file offset. The engine reads the compressed chunk containing it through the chunk cache or the OS page cache, decompresses the whole chunk, and hands rows to the merge. A large compression chunk length makes every small read pay to decompress the full chunk, which is why read-heavy tables with small rows often benefit from a smaller chunk_length_in_kb.

Merging sources and the price of tombstones

The merge iterator walks all sources in clustering order and, for each cell, keeps the version with the highest write timestamp. Deletions are just more data: a tombstone carries a timestamp and shadows anything older, whether it marks a cell, a row, a clustering range or the whole partition. The merge has to read the tombstones to know what to suppress, so a query that returns ten live rows may scan a hundred thousand dead ones first. Cassandra counts them: crossing tombstone_warn_threshold (1,000 by default) logs a warning and crossing tombstone_failure_threshold (100,000) aborts the query.

This is why the data model, not the hardware, sets most read latency. A queue table that deletes from the head and reads from the head forces every read to skip all the deletions since the last compaction purged them, which cannot happen until gc_grace_seconds has passed. A partition written across months under size-tiered compaction is spread across many SSTables, so one read touches many files. The number to watch is SSTables per read, which nodetool tablehistograms reports as a percentile distribution per table. Single digits at p99 is healthy; double digits means compaction strategy or partition design needs attention, and no amount of cache will hide it for long.

Reconciling replicas and blocking read repair

Back on the coordinator, the common case is cheap: the data response and digest agree, and the client gets the data. When they disagree, the coordinator issues a second round of full data requests to the replicas involved, merges them cell by cell with the same timestamp rule, and computes per-replica mutations that would bring each stale replica up to date.

Since Cassandra 4.0 the table option read_repair defaults to 'BLOCKING': the coordinator sends those repair mutations and waits for enough of them to be acknowledged before returning, so a value read at quorum cannot be followed by an older value at quorum. That is monotonic quorum reads, and it is why a mismatch shows up as extra latency on that one request rather than as a background task. The old read_repair_chance and dclocal_read_repair_chance options, which repaired a random fraction of reads in the background, were removed in 4.0. Setting read_repair = 'NONE' skips the write-back and gives up monotonic reads. The Cassandra documentation suggests it where multi-row writes to a partition must stay atomic, because blocking read repair can copy part of such a write to a stale replica. Read repair only fixes data you happen to read, so it never replaces regular anti-entropy repair.

Worked example: tracing a slow cart read

A checkout service reads a cart with SELECT * FROM cart_items WHERE cart_id = ?. p50 is 2 ms, p99 is 140 ms, and the team wants to know why. The first step is to stop guessing which stage is slow and trace a sample of real requests from the driver:

from cassandra import ConsistencyLevel
from cassandra.cluster import Cluster, ExecutionProfile, EXEC_PROFILE_DEFAULT
from cassandra.policies import (TokenAwarePolicy, DCAwareRoundRobinPolicy,
                                ConstantSpeculativeExecutionPolicy)

profile = ExecutionProfile(
    load_balancing_policy=TokenAwarePolicy(DCAwareRoundRobinPolicy(local_dc="dc1")),
    consistency_level=ConsistencyLevel.LOCAL_QUORUM,
    # client-side speculation: only applies to statements marked idempotent
    speculative_execution_policy=ConstantSpeculativeExecutionPolicy(delay=0.05, max_attempts=1),
    request_timeout=2.0,
)
cluster = Cluster(["10.0.1.11"], execution_profiles={EXEC_PROFILE_DEFAULT: profile})
session = cluster.connect("shop")

stmt = session.prepare("SELECT * FROM cart_items WHERE cart_id = ?")
stmt.is_idempotent = True

rs = session.execute(stmt, [cart_id], trace=True)   # sample 1 in 1000 in production, not every call
for ev in rs.get_query_trace().events:
    print(f"{ev.source_elapsed.total_seconds()*1000:8.2f} ms  {ev.source}  {ev.description}")

The trace lists events from the coordinator and each replica with elapsed time. In this case the slow traces showed the digest round trip completing in under 3 ms, but one replica's events read like Merged data from memtables and 14 sstables followed by Read 12 live rows and 9450 tombstone cells. nodetool tablehistograms shop cart_items confirmed a p99 of 14 SSTables per read, and the table was using size-tiered compaction with carts updated for weeks and items deleted on checkout.

Three changes followed, one per cause. Abandoned carts got a TTL so they expire instead of accumulating deletions. The table moved to leveled compaction (unified compaction in 5.0 can be configured to behave similarly), which bounds how many files hold one partition. And item removal changed from deleting rows one at a time to rewriting the cart as a single row-range deletion. After compaction caught up, SSTables per read fell to 2 at p99 and the query's p99 dropped to 9 ms. No cache or hardware changed.

Failure modes

The read path fails in a handful of recognisable ways, each with its own signature:

  • Read timeouts with a healthy cluster. ReadTimeoutException reports how many replicas responded versus how many were required and whether data was present. One-of-two with data present usually means one slow replica; check its disk latency and GC, and confirm speculative retry is not 'NONE'.
  • TombstoneOverwhelmingException. The failure threshold was crossed. Raising the threshold only moves the cliff; fix the access pattern or use TTLs and range deletes.
  • Latency spikes after a node outage. A replica that missed writes now mismatches on every read, so each request pays a second round trip plus blocking repair until anti-entropy repair runs.
  • Wide partitions. A partition of hundreds of megabytes makes compaction, streaming and the row index expensive and pins one replica set. Watch the large-partition warnings in the logs.
  • Cross-datacenter surprise. QUORUM instead of LOCAL_QUORUM counts replicas across all datacenters, so a read can block on a WAN round trip.

Trade-offs

KnobHelpsCosts
Token-aware routingOne fewer hop, local data readNone worth noting; enable it
speculative_retry (99p, fixed ms)Tail from a single slow replicaExtra replica reads under stress
Lower bloom_filter_fp_chanceFewer wasted SSTable seeksOff-heap memory per SSTable
Leveled or tuned unified compactionFew SSTables per readMore compaction I/O on writes
Smaller chunk_length_in_kbLess decompression per small readLower compression ratio, more metadata
read_repair = 'NONE'No write-back on mismatchLoses monotonic quorum reads
Row cacheHot small partitions served from memoryHeap or off-heap use, invalidation on writes

What to do next

  1. Confirm every client uses token-aware routing and a LOCAL consistency level in multi-datacenter clusters.
  2. Sample query traces for your slowest statements and classify each slow trace as coordinator wait, SSTable fan-out or tombstone scanning.
  3. Run nodetool tablehistograms on your top read tables; anything above single-digit SSTables per read at p99 gets a compaction or data-model review.
  4. Grep logs for tombstone warnings and large-partition warnings, and fix the top offender before tuning anything.
  5. Check speculative_retry and read_repair per table and write down why each is set the way it is.
  6. Make sure anti-entropy repair runs within gc_grace_seconds, so read repair is a safety net and not your consistency mechanism.

To keep going, read how Cassandra Bloom filters are sized, speculative retry in depth, read repair, tombstones and their cost and query tracing.

Key takeaway: A Cassandra read is three jobs: the coordinator sends one data request and enough digest requests to satisfy the consistency level, each replica narrows its SSTables with metadata, Bloom filters and an index before merging every source by timestamp, and the coordinator repairs any mismatch before it answers. Trace slow reads to find which job is slow, and remember that SSTables per read and tombstones scanned, both set by data model and compaction, decide most latency.