Chapter 27 · Observability, Performance, Capacity, Guardrails, Upgrades, Multi-DC Resilience, and Capstone
Capacity Planning for Partitions, SSTables, Compaction, Repair, Heap / Off-Heap, Disk, Network, and Failure Headroom
Turn measured partitions/SSTables/maintenance/memory/disk/network behavior into failure-aware capacity and N-1 headroom.
Learning outcomes
AtlasMart's three nodes average only 45% CPU and 55% disk usage, so a spreadsheet labels the cluster “half empty.” During repair plus one-node replacement, disk fills, compaction stalls and p99 explodes. Capacity is not normal-load utilization; it is the ability to meet the SLO through background maintenance and credible failure states.
Translate partition rows/bytes/cardinality and retention into per-node live-data/SSTable expectations.
Budget compaction, repair, snapshots/backup, streaming and tombstone overhead rather than multiplying logical bytes by RF only.
Treat heap, off-heap/page cache, disk and network as interacting constraints with measured saturation points.
Define N-1 and maintenance headroom using failure drills instead of a universal utilization percentage.
Produce a capacity sheet with demand, amplification, headroom, SLO and trigger-to-scale evidence.
The mandatory single-datacenter labs use Apache Cassandra
5.0.9 in the pinned
cassandra:5.0.9 image, Java 17 inside the
official image, cqlsh/nodetool from
that same image, and Apache Cassandra Java Driver
4.19.3 where client behavior matters. Use Docker
network atlasmart-cassandra-capstone, cluster
atlasmart-capstone, nodes
atlasmart-cap-1..3, datacenter dc1,
racks rack1..rack3, 16 virtual nodes (vnodes) per
node, NetworkTopologyStrategy, replication factor
(RF) 3, and application reads/writes at
LOCAL_QUORUM unless the lesson deliberately
changes consistency level (CL). Tables explicitly use
UnifiedCompactionStrategy (UCS),
default_time_to_live=0, and
gc_grace_seconds=864000.
For repeatability, mandatory Lessons 1–3 keep authentication/client TLS/internode TLS disabled only inside this isolated Docker network; no Cassandra or JMX port is published on the host, and JMX remains local-only inside containers. Chapter 26 remains the production security baseline. Lesson 4 creates separate disposable upgrade/multi-DC clusters, and Lesson 5 makes the security state an explicit capstone acceptance gate. A full RF=3 two-DC game day requires six Cassandra containers and roughly 10–14 GiB of available RAM depending on container limits/JVM ergonomics; the lesson also provides a reduced four-node RF=2 simulation for constrained laptops and labels the semantic difference.
Exact token values, latencies, GC pauses, SSTable counts, compaction/repair bytes, disk throughput, network rates, failure-detection timing, guardrail messages and benchmark throughput are runtime evidence. The lesson never treats example numbers as results from your machine. Before every benchmark/failure drill capture host CPU/RAM/disk, Docker limits, Cassandra/Java/driver versions, topology, RF/CL, schema/compaction, dataset shape, concurrency, warmup and security state.
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.
Terms used across the capstone
A coordinator is the Cassandra node handling
one client request; a replica stores a copy of
the requested partition according to the keyspace replication
strategy. A partition is the rows sharing a
partition key; the partitioner hashes that key to a
token, and virtual nodes
(vnodes) give each physical node multiple token
ranges. A datacenter (DC) and
rack model failure/locality domains.
RF (replication factor) is the number of
replicas per DC configured by
NetworkTopologyStrategy; a
CL (consistency level) controls how many
appropriately scoped replica responses are required.
An SSTable (Sorted String Table) is an immutable on-disk data-file set. Compaction rewrites SSTables to merge versions/tombstones and manage read/space amplification. Repair is anti-entropy comparison/streaming between replicas. heap is Java-managed memory; off-heap covers native/direct structures outside the Java heap; the operating-system page cache is another memory consumer. GC (garbage collection) reclaims Java heap objects and can introduce pauses.
A histogram is a distribution, not an average. p50/p95/p99 are percentiles: p99 means 99% of observed values are at or below that value. tail latency is the high-percentile response time users often feel during saturation/failure. throughput is completed work per time; error rate is failed work divided by attempted work. A Service Level Objective (SLO) is a target such as p99 latency/availability; RPO (Recovery Point Objective) is tolerated data-loss time; RTO (Recovery Time Objective) is tolerated service-recovery time.
JMX (Java Management Extensions) exposes
node-local Cassandra metrics/management operations;
nodetool is itself a JMX client. An
exported metric is forwarded to an external
time-series system; Cassandra metrics are node-local until an
operator aggregates them. A guardrail warns or
rejects dangerous schema/query/operational patterns. A
rolling upgrade changes one node at a time
while the cluster remains available.
schema agreement means nodes report the same
current schema version. A driver is the client
library implementing native protocol, topology discovery, load
balancing, timeouts, retries, idempotency and speculative
execution. SAI expands to Storage-Attached
Indexing; a vector index supports approximate
nearest-neighbor retrieval and has separate
memory/disk/build/recall costs.
1. Capacity starts at the partition, not at the server
Cassandra distributes partitions, not individual cells
independently. A bounded query-first partition lets you estimate
worst-case rows/bytes and reason about one read's work. Average
partition size is insufficient when a hot tenant or retention
bug creates a long tail. Use
tablehistograms partition-size/cell-count
percentiles and toppartitions activity sampling
together; the largest partition and hottest partition may be
different.
docker exec atlasmart-cap-1 nodetool tablehistograms atlasmart_capstone orders_by_customer_monthdocker exec atlasmart-cap-1 nodetool tablestats -H atlasmart_capstone.orders_by_customer_monthdocker exec atlasmart-cap-1 nodetool toppartitions atlasmart_capstone orders_by_customer_month 1000 10 --capacity 256 || true
| Layer | Capacity question | Evidence |
|---|---|---|
| partition | p95/max rows/bytes? bounded by month/bucket? | tablehistograms + business distribution |
| replicas | logical bytes × RF / topology balance | NTS/RF, nodetool status/tablestats |
| SSTables | files/read + total files/bytes | tablehistograms/tablestats |
| compaction | temporary rewrite bytes + throughput backlog | compactionstats/history + disk writes |
| repair | validation/anticompaction/stream bytes + unrepaired state | repair metrics/netstats/disk |
| backup/snapshot | hard-link retention + copied bytes | listsnapshots + backup catalog |
| SAI/vector | index bytes/build/query memory/recall | system_views.indexes + SAI/vector metrics |
2. Build an explicit capacity worksheet with ranges, not fake precision
# Fill these from workload + nodetool/exported metrics, not folklore.logical_live_gib = 600.0 # total application live data, before RFrf = 3nodes = 6index_ratio = 0.12 # observed SAI/other index bytes / data bytesmetadata_ratio = 0.08 # observed compression/index/metadata overhead range inputsnapshots_gib_per_node = 20.0 # retained local true-size evidencetemp_rewrite_ratio = 0.45 # observed/planned compaction/repair temporary headroomfailure_margin = 0.30 # policy from N-1 drill, not a Cassandra constantreplicated_cluster = logical_live_gib * rfsteady_per_node = replicated_cluster / nodessteady_with_index = steady_per_node * (1 + index_ratio + metadata_ratio)temp_headroom = steady_with_index * temp_rewrite_ratiopolicy_headroom = steady_with_index * failure_marginrequired_per_node = steady_with_index + snapshots_gib_per_node + temp_headroom + policy_headroomprint({ 'replicated_cluster_GiB': round(replicated_cluster,1), 'steady_per_node_GiB': round(steady_with_index,1), 'temporary_rewrite_headroom_GiB': round(temp_headroom,1), 'failure_policy_headroom_GiB': round(policy_headroom,1), 'planning_disk_per_node_GiB': round(required_per_node,1),})
The ratios above are intentionally inputs. Cassandra does not promise one compaction/index/repair amplification factor across workloads. UCS behavior, tombstones, partition overlap, SAI/vector indexes, compression, snapshots, repaired-state transitions and failure topology change the result. Store low/expected/high cases and validate them under load.
That free space may be simultaneously needed by compaction output, repair anti-compaction/streaming, snapshots, bootstrap/replacement and growth. The safe capacity question is whether the node can complete the largest credible overlapping maintenance/failure scenario while meeting the SLO and avoiding disk exhaustion.
3. Heap is not the whole memory budget
Java heap contains Cassandra objects/caches/memtables depending on implementation, but Cassandra also uses off-heap/direct/native memory and relies heavily on the operating-system page cache for SSTable reads. Giving the JVM every available byte can reduce page-cache effectiveness and increase host/container pressure. Conversely, too-small heap can raise allocation/GC frequency. Measure GC pause/time, heap occupancy after GC, off-heap/direct usage where exported, RSS and page-cache/host memory under representative concurrency.
docker exec atlasmart-cap-1 nodetool infodocker exec atlasmart-cap-1 nodetool gcstatsdocker exec atlasmart-cap-1 nodetool tpstatsdocker exec atlasmart-cap-1 nodetool compactionstatsdocker stats --no-stream atlasmart-cap-1 atlasmart-cap-2 atlasmart-cap-3# Linux host/container details vary by runtime; record what is actually exposed.docker exec atlasmart-cap-1 sh -lc 'cat /proc/meminfo | head -20'docker exec atlasmart-cap-1 sh -lc 'cat /sys/fs/cgroup/memory.max 2>/dev/null || true'
Boundary case: a benchmark that warms a tiny dataset entirely into page cache can make disk irrelevant; a production working set larger than cache can show very different p99. Record cold-ish/warm steady state separately.
4. Disk and network must cover foreground + background + failure traffic
docker exec atlasmart-cap-1 nodetool compactionstatsdocker exec atlasmart-cap-1 nodetool netstatsdocker exec atlasmart-cap-1 nodetool listsnapshotsdocker exec atlasmart-cap-1 nodetool tablestats -H atlasmart_capstone.orders_by_customer_monthdocker exec atlasmart-cap-1 sh -lc 'df -h /var/lib/cassandra; du -sh /var/lib/cassandra/data /var/lib/cassandra/commitlog /var/lib/cassandra/hints 2>/dev/null'
Repair and bootstrap/replacement are network consumers and can also trigger disk reads/writes/compaction. Cross-DC repair consumes WAN bandwidth. Capacity plans should reserve bandwidth for the recovery objective, not merely application traffic. If a 4 TiB node can stream only 50 MiB/s after application traffic and throttles, a full replacement is on the order of many hours before validation/compaction—not a 30-minute RTO.
5. N-1 headroom is an SLO experiment
In the three-node RF=3 teaching topology, every node owns a replica of every partition, so losing one node immediately reduces replica availability but does not redistribute new ownership until replacement/topology action. That topology is convenient but poor for extrapolating per-node storage balance. For capacity, repeat the same workload on a larger topology (for example six nodes RF=3) and measure application + repair/replacement behavior with one node unavailable.
# Baseline evidence first.docker exec atlasmart-cap-1 nodetool proxyhistogramsdocker stats --no-stream atlasmart-cap-1 atlasmart-cap-2 atlasmart-cap-3# Failure injection blast radius: this disposable cluster only.docker pause atlasmart-cap-3docker exec atlasmart-cap-1 nodetool status# Re-run the SAME measured workload profile at LOCAL_QUORUM and capture p95/p99/errors.docker exec atlasmart-cap-1 cassandra-stress user \ profile=/tmp/atlasmart-stress.yaml duration=60s \ "ops(insert=3,latest_orders=7)" no-warmup cl=LOCAL_QUORUM \ -node atlasmart-cap-1,atlasmart-cap-2 -rate threads=16docker unpause atlasmart-cap-3docker exec atlasmart-cap-1 nodetool statusdocker exec atlasmart-cap-1 nodetool repair --full atlasmart_capstone orders_by_customer_month
A passing N-1 capacity test means the defined workload/SLO survived that scenario with acceptable resource/recovery headroom. It does not prove rack/DC disaster, simultaneous repair+backup, a second failure, or future growth. Define which overlaps are credible and test them.
6. Capacity trigger sheet
| Trigger | Scale/redesign action | Evidence before action |
|---|---|---|
| partition p95/max violates model | bucket/shard/query-table redesign | partition histogram + query contract |
| SSTables/read/tail rises with compaction debt | fix throughput/headroom/model; review compaction | histograms + compactionstats |
| repair cannot finish within policy window | add capacity/bandwidth or split ranges/schedule | repair duration/bytes + SLO impact |
| heap/GC tail unacceptable before CPU saturation | profile allocation/cache; right-size JVM/container | GC + RSS/off-heap + workload |
| disk temp headroom insufficient | add disk/nodes, reduce snapshots/change workflow | high-water scenario worksheet |
| network limits replacement/RTO | add bandwidth/nodes/change streaming plan | netstats + measured transfer |
| N-1 p99/error SLO fails | scale before growth or reduce per-node responsibility | same workload under failure |
Check your understanding
- Why is logical data × RF / nodes not a complete disk plan?
- Why can increasing Java heap reduce read performance?
- What does an N-1 test prove?
- Why should network capacity include repair/replacement?
- What capacity input should come before node sizing?
Review the answers
1. It omits indexes/metadata, compaction/repair temporary writes, snapshots/backups, streaming, imbalance, growth and failure headroom.
2. It may steal memory from OS page cache/off-heap and can change GC behavior; the optimum is workload/host dependent.
3. Only that the tested workload/SLO/resource scenario survived the specified single failure and recovery sequence.
4. Recovery streams can dominate bandwidth and directly determine convergence/replacement/RTO.
5. Bounded partition/query/retention distributions and representative workload/SLO requirements.
Production judgment
Do not promote one local run into a universal Cassandra tuning rule. Record workload fit and non-goals; partition cardinality/rows/bytes and retention; read/write mix and p50/p95/p99/max latency; RF/CL and coordinator/replica failure behavior; JVM heap/GC, off-heap/page-cache, disk capacity/latency/IOPS/throughput and network bandwidth/packet loss; SSTable/read amplification, compaction backlog, tombstones and repair state; SAI/vector build/query/write and recall costs where used; authentication/authorization/TLS/JMX/secret/tenant boundaries; driver local-DC routing, timeout/retry/idempotency/speculation behavior; metrics/log/tracing coverage; backup RPO/RTO/restore proof; and the operator skill/runbook needed to perform repair, topology, upgrade and recovery safely.
Managed Cassandra services may hide disks, JMX, repair, backup, upgrade sequencing or metric names, and their quotas/cost model can change capacity decisions. Translate the same evidence questions into provider-native signals; do not assume the provider removes application data-model, driver, consistency, SLO, security, migration or rollback responsibility. Lesson 3 turns the risky patterns discovered by these measurements into guardrails and performance controls that warn/reject before an incident.
Summary and next step
This lesson’s concepts, evidence path, failure boundaries, and production judgment should now be explicit enough to verify rather than assume. Re-run the check-your-understanding prompts and preserve any lab evidence you need before changing or cleaning up the environment.
Next, continue to Guardrails and Performance Controls: Prevent Dangerous Schema/Query Patterns Before They Become Incidents.
Authoritative references
Re-check these version-sensitive sources before a real upgrade, capacity commitment, security change or game day.
- Apache Cassandra 5.0 release/download baseline
- Monitoring metrics
- Troubleshooting with nodetool histograms
- nodetool tablestats
- nodetool tablehistograms
- nodetool proxyhistograms
- cassandra-stress user mode
- cassandra.yaml including guardrails/storage compatibility
- nodetool setguardrailsconfig
- Repair operations
- Backups
- Security
- Java Driver core documentation