Cluster averages are often where a hot partition goes to hide.
Per-Partition/Shard Metrics, Replication Lag, Queue Depth, Compaction, Cache, and Storage Signals
Diagnose hidden hot shards and background debt by correlating per-partition traffic with queueing, replica lag, compaction, cache, tombstone, disk, network, and repair signals.
Explain why cluster averages can conceal a hot shard, lagging replica, overloaded queue, or compaction/tombstone problem.
Correlate per-shard traffic with replication lag, pending work, cache behavior, disk latency/space, network, runtime, repair, and rebalance state.
Diagnose an AtlasMart shard whose hot-key traffic couples with queueing, lag, compaction debt, cache misses, disk tail, and tombstone work.
Choose mitigation from the observed mechanism and verify the change instead of treating rebalancing or tuning as automatically successful.
1. The average cluster is not a machine
AtlasMart hashes most keys evenly, but a celebrity product, large tenant, monotonic time bucket, or badly chosen partition key can concentrate traffic. A cluster average of 53 queued requests is ambiguous: perhaps every shard has a small queue, or one shard has 190 while the others have single digits. The second case has a very different failure mode and remediation.
Always preserve the routing dimension: shard/partition/table/node/replica, operation class, and where practical the bounded hot-key class. A shard is a unit of routing/ownership; a replica is a copy of that shard. Their metrics answer different questions. Replication lag indicates copy freshness; queue depth indicates waiting work; compaction backlog indicates immutable-file maintenance debt; cache hit rate can explain extra storage reads; disk latency/space can expose I/O saturation; tombstone or obsolete-version work can multiply read amplification; repair/rebalance changes load while restoring redundancy.
2. Build an evidence chain from route to storage
Suppose shard s3 receives 60% of requests. Its coordinator queue rises, followers fall behind, cache hit rate collapses, SSTables accumulate, and disk p99 rises. Those observations form a plausible mechanism chain: concentrated reads/writes create background work; background work competes for disk/CPU; service time rises; the queue grows; lag increases; client tail latency follows. But correlation is not proof. A trace, storage latency, compaction history, or controlled workload split is needed to validate which edge in the chain is causal.
| Layer | Examples of per-scope evidence | Failure question |
|---|---|---|
| Routing | requests/shard, key distribution, fan-out, coordinator ownership | Is traffic concentrated or scatter-gathering? |
| Replication | lag, missing acknowledgements, hints/backlog | Is one copy stale or slow? |
| Queues/runtime | pending tasks, pool wait, GC pauses, thread saturation | Where is work waiting? |
| Storage engine | pending compactions, SSTable count, tombstones, flush/WAL waits | Is maintenance debt amplifying reads/writes? |
| Cache/disk/network | hit/miss/evictions, disk p99/space, retransmits/bytes | Is the working set or I/O path saturated? |
| Repair/rebalance | streaming bytes, repair progress, ownership movement | Is recovery/maintenance consuming headroom? |
A lower cache hit rate may be harmless if the cache is intentionally small and disk latency is healthy. A high SSTable count may be expected for a strategy/workload. Interpret metrics against workload shape, product semantics, and client symptoms.
3. Deliberately wrong approach: average all shards and tune the cluster
The lab has three healthy shards and one overloaded shard. Averaging cache hit rate to 0.76 or disk p99 to 14.25 ms obscures s3 at 0.34 cache hit and 41 ms disk p99. A blanket cluster-wide knob change could spend resources everywhere while leaving the hot routing key intact.
4. AtlasMart lab: expose a hidden hot shard
Python 3.13+ standard library only. The “rebalanced” values are explicit planning arithmetic, not measured results from Cassandra or another engine.
shards = {
"s0": {"traffic": 8200, "queue": 8, "lag_s": 0.8, "pending_compactions": 2, "cache_hit": .91, "disk_p99_ms": 5, "tombstones_per_read": 12},
"s1": {"traffic": 7600, "queue": 7, "lag_s": 1.1, "pending_compactions": 3, "cache_hit": .89, "disk_p99_ms": 6, "tombstones_per_read": 18},
"s2": {"traffic": 7900, "queue": 9, "lag_s": 0.9, "pending_compactions": 2, "cache_hit": .90, "disk_p99_ms": 5, "tombstones_per_read": 15},
"s3": {"traffic": 36500, "queue": 190, "lag_s": 43.0, "pending_compactions": 61, "cache_hit": .34, "disk_p99_ms": 41, "tombstones_per_read": 7800},
}
def avg(field):
return sum(v[field] for v in shards.values()) / len(shards)
print("CLUSTER AVERAGES")
for field in ("queue","lag_s","pending_compactions","cache_hit","disk_p99_ms"):
print(field, round(avg(field),2))
print("\nPER-SHARD EVIDENCE")
for name,v in shards.items():
print(name, v)
hot = max(shards, key=lambda s: shards[s]["traffic"])
print("\nhot shard:", hot, "traffic_share=", f"{shards[hot]['traffic']/sum(v['traffic'] for v in shards.values()):.1%}")
print("diagnosis: shard s3 couples hot traffic with queueing, lag, compaction debt, cache misses, disk tail, and tombstone work")
print("\nMODELED REBALANCE / KEY-SPLIT RESULT")
# Move half of s3's celebrity-key work to a new bucket; this is planning arithmetic, not a benchmark.
rebalanced = dict(shards["s3"])
rebalanced["traffic"] = 19000
rebalanced["queue"] = 42
rebalanced["lag_s"] = 7.0
rebalanced["pending_compactions"] = 18
rebalanced["cache_hit"] = .71
rebalanced["disk_p99_ms"] = 13
rebalanced["tombstones_per_read"] = 1200
print("s3 before:", shards["s3"])
print("s3 after :", rebalanced)
print("verify with production metrics; a modeled improvement is not proof of an actual fix")
Shard s3 carries 60.6% of traffic and simultaneously shows queue=190, lag=43 s, 61 pending compactions, 0.34 cache hit, 41 ms disk p99, and 7,800 tombstones/read in the toy model. Splitting the hot work reduces modeled pressure, but the lesson requires production verification because arithmetic is not a benchmark.
5. Production judgment: mitigate the mechanism, preserve failure headroom
If the root cause is a hot key, adding random capacity may not help if routing still pins the key to one owner. Options include changing the partition/bucket key, precomputing a read projection, separating heavy tenants, request coalescing, or—in carefully defined read-only cases—replicating a hot object. Each changes consistency, fan-out, repair, or cost. If background maintenance is the bottleneck, verify compaction/repair throughput and temporary disk headroom rather than simply increasing foreground concurrency.
Monitor the maximum and high-percentile shard values in addition to averages. Record topology generation during incidents so “shard s3” means the same ownership epoch. The next lesson turns these observations into a benchmark contract that reproduces realistic skew instead of feeding the cluster uniformly.
Check your understanding
- Why are cluster averages dangerous for partitioned databases?
- What does pending compaction work tell you?
- Why might adding nodes fail to fix a celebrity key?
- What should be recorded during rebalance or repair?
- Why is modeled post-fix improvement not proof?
Review the answers
1. They can combine healthy and pathological shards, hiding the maximum queue, lag, disk tail, hot-key share, or maintenance debt that dominates user experience.
2. It indicates storage-maintenance debt is queued; its operational significance depends on rate, disk headroom, strategy, read/write amplification, and whether the backlog grows.
3. If the routing key remains indivisible, that key can still map to one shard/replica set; cluster capacity grows without removing the local bottleneck.
4. Ownership/topology generation, streaming/repair progress, client latency/errors, queue/storage metrics, and remaining redundancy/headroom.
5. Real systems include caches, disks, scheduling, network, replication, concurrency, and product behavior absent from a simple arithmetic model; rerun representative measurements.
References
Foundational claims use primary research or current official documentation where practical. Product references are implementation anchors only; the mandatory labs are vendor-neutral.
- Apache Cassandra 5.0 — Monitoring — Official per-table/node cache, pending compaction, dropped-message, commit-log, streaming, and repair metrics.
- Apache Cassandra 5.0 — Compaction Overview — Official description of SSTables, tombstones, compaction work, and space/performance consequences.
- Apache Cassandra 5.0 — Troubleshooting with nodetool — Official operational interpretation of pending thread-pool work, compactions, drops, and latency.
- PostgreSQL 18 — Monitoring Database Activity — Current official contrast showing node/database activity, replication, I/O, WAL, and wait-oriented signals in another engine family.