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.

Advanced125–165 minutesShard-diagnostics labPython 3.13+ · standard libraryVendor-neutral · free/local mandatory pathLast reviewed: August 2026
01

Explain why cluster averages can conceal a hot shard, lagging replica, overloaded queue, or compaction/tombstone problem.

02

Correlate per-shard traffic with replication lag, pending work, cache behavior, disk latency/space, network, runtime, repair, and rebalance state.

03

Diagnose an AtlasMart shard whose hot-key traffic couples with queueing, lag, compaction debt, cache misses, disk tail, and tombstone work.

04

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?
Do not “fix” one metric in isolation

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

Mandatory lab environment

Python 3.13+ standard library only. The “rebalanced” values are explicit planning arithmetic, not measured results from Cassandra or another engine.

python · AtlasMart deterministic simulation
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")
Expected evidence

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

  1. Why are cluster averages dangerous for partitioned databases?
  2. What does pending compaction work tell you?
  3. Why might adding nodes fail to fix a celebrity key?
  4. What should be recorded during rebalance or repair?
  5. 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.

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.