Every time a Cassandra cluster changes shape, data has to move between nodes. A node joins and must receive the ranges it now owns. A node leaves and must hand its ranges to whoever inherits them. A replacement node must rebuild a dead peer's data. A new datacenter must be filled from an existing one. Repair finds replicas that disagree and ships the missing rows. All of these use the same machinery: streaming, the subsystem that moves SSTable data from one node's disk to another's.
Operators usually learn each operation as a separate ritual, which is why the same mistakes keep recurring: a throttle set in the wrong unit, a session that dies at 97% and starts again from zero, a receiving disk that fills before compaction catches up. This article treats streaming as one subsystem with one lifecycle, one set of throttles and one set of failure modes.
Settings and commands below follow the current Cassandra documentation (the 5.0 configuration reference and nodetool pages) unless a different version is named. Older releases used different key names for the same settings, so check your own cassandra.yaml before copying values. Token allocation and the details of the bootstrap plan are covered in the bootstrap architecture article; this one is about running and debugging streaming in general.
One subsystem behind every operation
The operations differ in three ways: which node starts the work, which nodes send data, and which token ranges move. If you know those three things, you can predict an operation's load, duration and failure behaviour.
| Operation | Started on | Sources | What moves |
|---|---|---|---|
| Bootstrap (new node joins) | Joining node | Current replicas of its new ranges | Every range the new node now owns, in full |
| Host replacement | Replacement node, started with -Dcassandra.replace_address_first_boot=<dead ip> | Surviving replicas | The dead node's ranges, in full |
nodetool decommission | Leaving node | The leaving node itself | Its ranges, to each range's next owner |
nodetool removenode <host id> | Any live node | Surviving replicas of the dead node's ranges | Ranges that lost a replica |
nodetool rebuild <src-dc> | The node being filled | Replicas in the named datacenter | Everything the node owns, typically a new datacenter |
| Repair | Coordinator of the repair | Replicas that disagree | Only the subranges whose Merkle trees differ |
Bulk load (sstableloader) | An external client | The loader process | Prepared SSTables, split by owner |
Two rows matter most. Decommission has one source, the leaving node, so that node's outbound throttle limits the whole operation. Bootstrap, rebuild and replacement have many sources, so total throughput is roughly the sum of several throttled senders. The receiver's disk, CPU and compaction then become the bottleneck. Repair moves small, scattered pieces instead (see the repair architecture article).
Anatomy of a stream session
A streaming operation starts as a stream plan on the node that begins it. The plan has an ID, an operation type (Bootstrap, Rebuild, Repair, Unbootstrap and so on) and one session per peer. Every session goes through the same three phases, shown below.
- Prepare. The two nodes agree on what will move: which keyspaces and tables, which token ranges, and in which direction. A session can send, receive or both. Repair sync sessions often move data both ways.
- Transfer. The sender reads the SSTables that hold the requested ranges and sends them file by file; the receiver acknowledges each one. If an entire SSTable falls inside the requested ranges, it can be sent whole: the data file and its index, filter and metadata components are copied without being decoded. Otherwise the sender reads only the matching partitions and the receiver writes them into new SSTables.
stream_entire_sstables(default true) enables the whole-file path. - Complete. Both sides confirm that everything arrived. On the receiver, the new SSTables for a table join the live data set only when that table's receive task finishes, so a half-transferred table is never readable.
While a session is otherwise idle, the two nodes exchange keep-alive messages every streaming_keep_alive_period (default 300s). They turn a silently dropped connection into a failure after a bounded wait, and a failed session can be retried. streaming_connections_per_host (default 1) controls how many connections a session opens to each peer.
There is one important exception on the receive side. Tables with materialized views (and, depending on version and settings, tables with change data capture) cannot simply have files dropped into place, because each incoming row must also update the view or the CDC log. Those streams are replayed through the normal write path, which is much slower and puts load on memtables and the commit log. If a rebuild is far slower than the data volume suggests, check for materialized views first.
Throttles and the arithmetic of a change window
Streaming has four outbound throttles. All four default to 24MiB/s in the 5.0 configuration, and each one is a per-node limit on data that node sends:
# cassandra.yaml (5.0 names; values shown are the defaults)
stream_entire_sstables: true
stream_throughput_outbound: 24MiB/s # partition-by-partition streaming
inter_dc_stream_throughput_outbound: 24MiB/s # same, to other datacenters
entire_sstable_stream_throughput_outbound: 24MiB/s # whole-SSTable streaming
entire_sstable_inter_dc_stream_throughput_outbound: 24MiB/s
streaming_keep_alive_period: 300s
streaming_connections_per_host: 1Notice that whole-SSTable transfers have their own pair of limits. If you raise only stream_throughput_outbound and most of your data travels as whole files, nothing gets faster. Both pairs can be changed at runtime, but the nodetool units are a trap:
nodetool setstreamthroughput 200 # bare value is in MEGABITS per second (~23.8 MiB/s)
nodetool setstreamthroughput -m 100 # -m: 100 MiB/s
nodetool setstreamthroughput -e -m 100 # -e: the entire-SSTable throttle instead
nodetool setstreamthroughput 0 # 0 disables throttling entirely
nodetool getstreamthroughput # confirm what is actually setSomeone who types setstreamthroughput 100 to mean 100 MiB/s has actually set 100 megabits, about 12 MB/s, which is roughly half the default. Runtime changes also disappear when the node restarts, so if the change must survive a restart, put it in cassandra.yaml as well.
Worked example. Suppose you are decommissioning a node that holds 1.2 TiB in a six-node cluster with RF=3. Decommission has a single source, so the leaving node's throttle sets the pace. 1.2 TiB is about 1,258,291 MiB. At the default 24 MiB/s that takes about 52,400 seconds, or 14.6 hours. At 100 MiB/s it takes about 12,600 seconds, or 3.5 hours, provided the disks and network can sustain it. If you leave the throttle in megabits at 100 (about 12.5 MB/s), the same job takes more than a day. A bootstrap pulling the same data from five peers at 24 MiB/s each has a 120 MiB/s ceiling, so the joining node's disk and compaction usually set the limit instead.
def stream_hours(data_tib, mib_per_s, sources=1, receiver_cap_mib_s=None):
"""Lower bound on wall time for a streaming operation."""
rate = mib_per_s * sources # each source is throttled independently
if receiver_cap_mib_s:
rate = min(rate, receiver_cap_mib_s) # disk/compaction limit on the receiver
return data_tib * 1024 * 1024 / rate / 3600
print(round(stream_hours(1.2, 24), 1)) # decommission, default: 14.6
print(round(stream_hours(1.2, 100), 1)) # decommission at -m 100: 3.5
print(round(stream_hours(1.2, 24, 5, 80), 1)) # bootstrap, 5 peers, receiver at 80 MiB/s: 4.4Treat these numbers as minimums. Large partitions, materialized views and a backed-up compaction queue on the receiver all add time.
Watching a stream
Use all three of these views; each shows something the others miss.
nodetool netstatslists active stream plans by operation and plan ID, then, for each peer, how many files and bytes it is sending or receiving and how many have arrived so far. Run it on both ends to see which side is stuck.- The streaming virtual table. Newer releases expose streaming state through a
system_views.streamingtable, sized and retained bystreaming_state_size(40MiB) andstreaming_state_expires(3d), withstreaming_stats_enabledas the switch. It is useful because it keeps recent sessions, including failed ones, after netstats has forgotten them. RunDESCRIBE TABLE system_views.streamingon your version to see which columns you have. - The system log. Each session logs when it starts, when it completes and, if it fails, the reason. Search the log for the plan ID that netstats shows. The first error for that plan is usually the cause.
# Watch a running operation from the receiving node, once a minute
watch -n 60 'nodetool netstats | grep -v "100%"'
# After the fact: what happened to recent sessions (5.0-style vtable)
cqlsh -e "SELECT * FROM system_views.streaming;"
# Correlate with the log by plan id
grep -n "<plan-id>" /var/log/cassandra/system.log | head -50Also watch the effects on the rest of the node, not just the bytes moving. Watch pending compactions on the receiver (nodetool compactionstats), free disk on the receiver, and p99 read latency on the sources.
Failure modes and recovery
Most streaming incidents fit one of these patterns.
- The session fails late and the work is lost. A network blip, a keep-alive timeout or a peer restart fails the session. Bootstrap can be continued with
nodetool bootstrap resume, which skips ranges that already finished. Its--forceflag overrides thecassandra.reset_bootstrap_progresssetting, and the docs mark that as dangerous. Rebuild has no resume, but you can narrow a re-run with-ksfor one keyspace,-tsfor specific token ranges, or-sfor specific source hosts, so you only re-stream what is missing. - One partition holds up the session. Partition-by-partition streaming cannot split a single partition. A multi-gigabyte partition is one indivisible unit of work, and a timeout while it is being sent sends you back to the start. The fix is in the data model.
- The receiver runs out of disk. Streamed data arrives before compaction has merged it with what the node already holds, so peak disk use can briefly exceed the node's final data size. On the old owners, data the cluster no longer needs there is only dropped when you run
nodetool cleanupafter the topology change. Plan for both kinds of headroom. - Compaction falls behind after repair. Repair streams many small pieces, and each becomes a small SSTable on the receiver. With leveled compaction, this lands as a pile of level-0 files and read latency rises until compaction catches up. Repair in smaller subranges and watch pending compactions (see compaction strategies).
- removenode loses its sources.
removenodestreams from the surviving replicas. If a range has no healthy replica left, it cannot be restored from anywhere.nodetool removenode statusshows progress, andremovenode forcefinishes the removal without completing the streams. After a forced removal, run a full repair, and accept that ranges with no surviving copy are gone. - Mixed versions. Finish a rolling upgrade before any topology change or repair.
- Encryption assumptions. The 4.0 streaming docs say zero-copy is disabled under internode encryption; on an encrypted cluster, confirm the fast path from throughput before counting on it.
A runbook for planned streaming
This pattern works for any planned operation; only step 3 changes.
#!/usr/bin/env bash
set -euo pipefail
NODE=10.0.3.17
# 0. Preconditions: cluster healthy, no upgrade in progress, no other topology change
if nodetool status | grep -E '^(DN|UL|UJ|UM) '; then echo "cluster not settled"; exit 1; fi
# 1. Headroom: the inheriting nodes need free disk for the incoming data plus compaction space
nodetool -h "$NODE" info | grep -i load
# 2. Raise throttles on the sending node for the window (runtime only)
nodetool -h "$NODE" setstreamthroughput -m 100
nodetool -h "$NODE" setstreamthroughput -e -m 100
# 3. The operation itself (blocks until done or failed)
nodetool -h "$NODE" decommission
# 4. Verify, then restore throttles on any node you changed and still run
nodetool statusChange one node at a time. Raise throttles in steps and stop when p99 latency on the sources moves. After a bootstrap or replacement completes, run nodetool cleanup on the nodes that gave up ranges, one at a time, to reclaim the space.
Trade-offs
Speed against live traffic. Higher throttles shorten the window spent at reduced redundancy but take disk and network from client requests.
Whole-file against partition streaming. Whole-SSTable transfer avoids deserialising and rewriting data, so it is much cheaper for both sides. It only applies when an SSTable sits entirely within the ranges being moved. That is common with leveled compaction and rare with large size-tiered files, so your compaction strategy affects how fast your topology changes are.
Rebuild against repair. Use rebuild to fill an empty node predictably; use repair, which builds Merkle trees and copies only differences, to keep filled nodes consistent.
Forcing against waiting. removenode force and bootstrap resume --force exist for emergencies. They trade completeness for time, and the only way to get the completeness back is a repair.
What to do next
- Open cassandra.yaml on one node and write down all four outbound throttles and
stream_entire_sstables. Check whether your release uses the unit-suffixed key names. - Run
nodetool getstreamthroughputon every node and compare with the file. Any differences are runtime changes someone made and forgot. - For your largest node, work out the decommission time at the default throttle and at the throttle you would actually use. Write both numbers into the runbook.
- Check whether your version has
system_views.streaming. If it does, add a query against it to your topology-change runbook. - List your tables that have materialized views or CDC, and expect those to stream through the write path.
- Find your largest partitions (
nodetool tablehistograms), because those are the units that will stall a stream. - Rehearse a failed bootstrap on staging: kill the network mid-stream, then practise
nodetool bootstrap resumeand a narrowednodetool rebuild -ks. - Read the SSTable format to understand which components travel in a whole-file transfer.