Chapter 11 · SSTables and On-Disk Storage Internals
Streaming SSTables During Bootstrap, Repair, Rebuild, and Topology Changes
Connect SSTable streaming to bootstrap, repair, rebuild, replacement, topology changes, network throttles, and disk headroom.
Learning outcomes
AtlasMart plans to add capacity and assumes “streaming just copies the current dataset once.” In reality bootstrap, rebuild, repair, replacement, and range movement can stream different token ranges, interact with compaction, create temporary disk/network pressure, and leave cleanup work behind. This lesson follows SSTables through those topology workflows without performing a risky full production-like rebalance.
Explain why bootstrap, rebuild, repair, replacement, decommission/move can stream replica data for different reasons.
Use nodetool netstats and repair preview to observe or estimate streaming without inventing transfer numbers.
Explain entire-SSTable/zero-copy streaming eligibility versus section/partition streaming and format/topology dependencies.
Connect streaming to disk headroom, network throttles, compaction, cleanup, hints, and post-topology repair requirements.
Use a small safe local exercise plus deterministic planning math instead of requiring an expensive multi-DC lab.
The mandatory labs continue the disposable AtlasMart
environment used by earlier chapters: pinned
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 11 uses keyspace
atlasmart_storage with
NetworkTopologyStrategy, replication factor (RF)
3, normally LOCAL_QUORUM, and tables that
explicitly use UnifiedCompactionStrategy (UCS). Cassandra's
current default SSTable format is BIG unless
sstable.selected_format is changed; Cassandra 5.0
also supports BTI trie-indexed SSTables. Authentication,
client TLS, internode TLS, and remote JMX stay disabled only
inside this isolated learning network. The Apache Cassandra
Java Driver 4.19.3 is optional; mandatory storage evidence
uses cqlsh, nodetool, Docker/Linux
filesystem tools, and Cassandra's bundled SSTable utilities.
Re-check nodetool version,
java -version, actual
cassandra.yaml, disk free space, and selected
SSTable format before interpreting output.
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.
Storage terms before touching the filesystem
An SSTable (Sorted String Table) is Cassandra's immutable on-disk representation produced by a memtable flush, compaction, streaming, bulk load, or related storage workflow. A component is one file belonging to an SSTable generation; component sets and names depend on SSTable format/version. A partition is the rows sharing a partition key and is sorted with other partitions by token order inside SSTable data; rows inside a partition follow clustering order. Compression chunks are independently compressed blocks of the Data component, letting Cassandra read/decompress only relevant chunks instead of the full file. Compaction reads SSTables and writes replacement SSTables, then retires old ones when safe. Streaming transfers replica data between nodes for bootstrap, rebuild, repair, replacement, and topology movement. Disk headroom is free capacity reserved not only for live data but for temporary overlap during these operations. Offline SSTable tools inspect or transform SSTables outside normal CQL/native-protocol execution; many explicitly require Cassandra to be stopped and therefore must never be pointed casually at live production paths.
1. Same transport family, different operational intent
| Operation | Why data streams | Key postcondition |
|---|---|---|
| Bootstrap/add node | new node acquires token ranges it will replicate | new node becomes UP; old owners may need cleanup after range movement |
| Rebuild | existing node reconstructs local replicas from other nodes, often from selected DC | validate rebuilt ranges/data; monitor source/target load |
| Repair | replicas compare common token ranges and stream differences | repaired state/convergence according to repair type |
| Replace dead node | replacement acquires ranges/data of dead endpoint | may still require repair depending on missed-write/hint window conditions |
| Move/decommission/remove | ownership changes and data is transferred/retained safely | validate new ownership, then cleanup where appropriate |
Do not equate “streaming completed” with “cluster is fully repaired.” Streaming is a transport mechanism used by several workflows. Each workflow has its own consistency, ownership, cleanup, and validation semantics.
2. Observe current stream state and configuration before creating traffic
docker exec atlasmart-cass-1 nodetool versiondocker exec atlasmart-cass-1 nodetool statusdocker exec atlasmart-cass-1 java -versiondocker exec atlasmart-cass-1 sh -lc "grep -n -A8 -B2 '^sstable:' /etc/cassandra/cassandra.yaml || true"docker exec atlasmart-cass-1 sh -lc "df -h /var/lib/cassandra && df -i /var/lib/cassandra"
CREATE KEYSPACE IF NOT EXISTS atlasmart_storageWITH replication = {'class':'NetworkTopologyStrategy','dc1':3};CREATE TABLE IF NOT EXISTS atlasmart_storage.orders_by_customer_day ( customer_id text, order_day date, order_time timestamp, order_id uuid, status text, total decimal, note text, PRIMARY KEY ((customer_id, order_day), order_time, order_id)) WITH CLUSTERING ORDER BY (order_time DESC, order_id ASC) AND compaction = {'class':'UnifiedCompactionStrategy'} AND compression = {'class':'LZ4Compressor','chunk_length_in_kb':'16'};CONSISTENCY LOCAL_QUORUM;INSERT INTO atlasmart_storage.orders_by_customer_day(customer_id,order_day,order_time,order_id,status,total,note)VALUES ('cust-42','2026-09-07','2026-09-07T18:00:00Z',00000000-0000-0000-0000-000000000001,'PAID',129.90,'first');INSERT INTO atlasmart_storage.orders_by_customer_day(customer_id,order_day,order_time,order_id,status,total,note)VALUES ('cust-42','2026-09-07','2026-09-07T18:01:00Z',00000000-0000-0000-0000-000000000002,'PACKING',89.50,'second');INSERT INTO atlasmart_storage.orders_by_customer_day(customer_id,order_day,order_time,order_id,status,total,note)VALUES ('cust-77','2026-09-07','2026-09-07T18:02:00Z',00000000-0000-0000-0000-000000000003,'CREATED',44.00,'third');DESCRIBE TABLE atlasmart_storage.orders_by_customer_day;SELECT * FROM atlasmart_storage.orders_by_customer_dayWHERE customer_id='cust-42' AND order_day='2026-09-07';
docker exec atlasmart-cass-1 nodetool netstats -Hdocker exec atlasmart-cass-1 sh -lc "grep -n -E 'stream_entire_sstables|stream_throughput|inter_dc_stream' /etc/cassandra/cassandra.yaml || true"docker exec atlasmart-cass-1 nodetool getstreamthroughput 2>/dev/null || truedocker exec atlasmart-cass-1 sh -lc "df -h /var/lib/cassandra"docker exec atlasmart-cass-1 nodetool compactionstats
If no topology/repair operation is active,
netstats should show no current sending/receiving
sessions. This is a baseline, not a guarantee that background
network traffic is zero. Current Cassandra configuration allows
entire-SSTable streaming when eligible; otherwise Cassandra can
stream selected ranges/sections and materialize appropriate
SSTables on the receiver. Eligibility depends on format,
ownership boundaries, operation, and configuration.
3. Safe preview: estimate repair streaming before transferring data
Rather than forcing a large bootstrap in a three-node laptop cluster, use repair preview to build Merkle trees and estimate the differences that a repair would stream without performing the stream. On a freshly consistent tiny fixture, the estimate may be zero. To make the mechanism visible without unsafe network/firewall manipulation, create a controlled stale replica exactly as in Chapter 10, then preview a full repair. Exact byte estimates depend on timing and replica participation.
docker exec atlasmart-cass-1 nodetool disablehandoffdocker exec atlasmart-cass-2 nodetool disablehandoffdocker pause atlasmart-cass-3docker exec atlasmart-cass-1 cqlsh -e "CONSISTENCY LOCAL_QUORUM; UPDATE atlasmart_storage.orders_by_customer_day SET note='missed-by-node3' WHERE customer_id='cust-42' AND order_day='2026-09-07' AND order_time='2026-09-07T18:00:00Z' AND order_id=00000000-0000-0000-0000-000000000001;"docker exec atlasmart-cass-1 nodetool enablehandoffdocker exec atlasmart-cass-2 nodetool enablehandoffdocker unpause atlasmart-cass-3docker exec atlasmart-cass-1 nodetool status
# Preview estimates streaming; it does not stream the differences.docker exec atlasmart-cass-1 nodetool repair --preview --full atlasmart_storage orders_by_customer_day# In a second terminal while an actual repair runs, watch netstats.docker exec atlasmart-cass-1 nodetool repair --full atlasmart_storage orders_by_customer_daydocker exec atlasmart-cass-1 nodetool netstats -H
Repair may complete before you capture netstats,
and container networking/storage do not represent production.
The evidence goal is to understand session/range/byte fields
and to record a zero/no-op outcome honestly when there is
nothing to stream.
4. Entire-SSTable streaming and disk-headroom coupling
Cassandra can stream eligible entire SSTables, avoiding serialization of individual partitions and reducing CPU/GC pressure. But the receiver still needs space for incoming data and later compaction. If range boundaries require only sections of an SSTable, or eligibility conditions are not met, data can be streamed more granularly. In either case, planning only network bandwidth is insufficient: source read I/O, receiver write I/O, compaction, repair/anti-compaction, and snapshots can overlap.
| Capacity question | Evidence to record |
|---|---|
| How much may move? | repair preview/topology ownership, node load, table sizes, range estimates |
| How fast may it move? | stream throughput settings, network capacity, source/target disk throughput, concurrent streams |
| Can receiver hold it? | filesystem free bytes/inodes plus compaction/snapshot/restore headroom |
| What else competes? | foreground reads/writes, compaction, repair, backup, GC/JVM and network |
| What happens after ownership change? | validate status/schema/data, then documented cleanup/repair depending on operation |
A dangerous approach is to set streaming throughput to unlimited because a maintenance window is short. That can starve foreground traffic or disks. Throughput values are environment-specific; measure source/target saturation and adjust in a controlled window with rollback.
Check your understanding
- Does “streaming complete” mean repair is complete for the whole cluster?
- What does nodetool repair --preview do?
- Why can entire-SSTable streaming be faster?
- Why can a bootstrap increase disk pressure on both old and new nodes?
- Why not disable stream throttling universally?
Review the answers
1. No. Streaming is a data-transfer mechanism used by multiple operations; repair scope and cluster-wide convergence require their own validation.
2. It builds/compares repair structures and estimates streaming for the requested repair without performing the stream.
3. Eligible complete SSTables can be transferred without serializing individual partitions, reducing CPU/GC overhead.
4. Sources read/retain data while the new node receives it; compaction/snapshots may overlap, and old owners keep out-of-range data until cleanup.
5. Network and disk saturation can damage foreground latency and cluster stability; tune from measured capacity and workload priorities.
Production judgment
SSTable files expose valuable operational evidence, but they are not an application contract. Production decisions must combine logical data shape with physical storage state: partition rows/bytes, mutation and TTL/delete rates, RF/CL, read/write p95/p99, SSTables per read, compaction strategy and backlog, repaired/unrepaired state, compression ratio, CPU/decompression cost, disk throughput/latency, temporary compaction/streaming space, snapshot/backup retention, topology changes, repair cadence, and restore objectives. Include JVM/GC, page cache/off-heap use, driver timeouts/retries/idempotency, network bandwidth, tenant isolation, encryption-at-rest expectations, filesystem/device behavior, and managed-service restrictions.
Do not plan disk capacity as “live dataset bytes × RF” only. Flush, compaction, repair, streaming, snapshots, incremental backups, anti-compaction, restore staging, and operational safety margins can temporarily retain additional SSTables. Do not manually delete or edit component files to recover space. If space is critical, first stop unsafe automation, measure ownership/snapshots/compaction/streaming, and choose a documented recovery path with rollback. Lesson 5 combines component size, SSTable count, compaction/read amplification and operation-specific temporary copies into a practical disk-headroom model.
Summary and next bridge
Streaming moves replica ownership or repairs differences; the SSTable transport path is only one part of the operation. Monitor sessions, disk/network pressure and follow-up cleanup/repair separately. The final lesson turns those physical files and transfers into a disk-headroom/capacity model.
Authoritative references
These are version-sensitive sources of truth. Re-check them when regenerating the lesson because SSTable formats, utilities, defaults, and topology procedures evolve.
- Apache Cassandra downloads / current GA baseline
- Apache Cassandra storage engine and SSTable components
- Cassandra 5.0 cassandra.yaml SSTable format and streaming settings
- Apache Cassandra compression guidance
- Apache Cassandra SSTable tools safety overview
- sstablemetadata
- sstabledump
- sstablepartitions
- nodetool netstats
- Repair and streaming differences
- Topology changes and netstats monitoring
- nodetool bootstrap