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.

OperationStarted onSourcesWhat moves
Bootstrap (new node joins)Joining nodeCurrent replicas of its new rangesEvery range the new node now owns, in full
Host replacementReplacement node, started with -Dcassandra.replace_address_first_boot=<dead ip>Surviving replicasThe dead node's ranges, in full
nodetool decommissionLeaving nodeThe leaving node itselfIts ranges, to each range's next owner
nodetool removenode <host id>Any live nodeSurviving replicas of the dead node's rangesRanges that lost a replica
nodetool rebuild <src-dc>The node being filledReplicas in the named datacenterEverything the node owns, typically a new datacenter
RepairCoordinator of the repairReplicas that disagreeOnly the subranges whose Merkle trees differ
Bulk load (sstableloader)An external clientThe loader processPrepared 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.

Coordinator nodebuilds the stream planPeer Asession 1Peer Bsession 2Peer Csession 3prepare: ranges, tablesprepareprepare1. Prepareagree what moves2. Transferfiles + per-file acks3. Completeboth sides confirmFailedsession abortederror or keep-alive timeoutReceiver: SSTables written to disk, joined to the live set when the table task completesthen normal compaction merges them with existing data; partial files from a failed session are discarded
One stream plan, one session per peer. Each session prepares, transfers and completes independently; any error or missed keep-alive fails that session.
  1. 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.
  2. 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.
  3. 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: 1

Notice 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 set

Someone 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.4

Treat 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 netstats lists 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.streaming table, sized and retained by streaming_state_size (40MiB) and streaming_state_expires (3d), with streaming_stats_enabled as the switch. It is useful because it keeps recent sessions, including failed ones, after netstats has forgotten them. Run DESCRIBE TABLE system_views.streaming on 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 -50

Also 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 --force flag overrides the cassandra.reset_bootstrap_progress setting, and the docs mark that as dangerous. Rebuild has no resume, but you can narrow a re-run with -ks for one keyspace, -ts for specific token ranges, or -s for 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 cleanup after 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. removenode streams from the surviving replicas. If a range has no healthy replica left, it cannot be restored from anywhere. nodetool removenode status shows progress, and removenode force finishes 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 status

Change 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

  1. 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.
  2. Run nodetool getstreamthroughput on every node and compare with the file. Any differences are runtime changes someone made and forgot.
  3. 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.
  4. Check whether your version has system_views.streaming. If it does, add a query against it to your topology-change runbook.
  5. List your tables that have materialized views or CDC, and expect those to stream through the write path.
  6. Find your largest partitions (nodetool tablehistograms), because those are the units that will stall a stream.
  7. Rehearse a failed bootstrap on staging: kill the network mid-stream, then practise nodetool bootstrap resume and a narrowed nodetool rebuild -ks.
  8. Read the SSTable format to understand which components travel in a whole-file transfer.
Key takeaway: Bootstrap, replacement, decommission, removenode, rebuild, repair and bulk loading all use one streaming subsystem: a plan of per-peer sessions that prepare, transfer and complete. Work out the number of sources, then calculate the window from the per-node throttles, remembering that whole-SSTable transfers have their own limits and that bare nodetool values are megabits. Watch netstats, the streaming virtual table and the log together, and keep resume, narrowed rebuilds and repair ready for when a session fails.