Cassandra's pitch is linear scalability: add nodes and capacity grows in proportion. The pitch is true, but the operation is not instant. Every topology change moves data across the network, rewrites SSTables and briefly changes who owns which tokens. Scaling down is riskier than scaling up, because it concentrates the same data onto fewer disks and can quietly cut redundancy below the replication factor you designed for.

This article is about the decision and the runbooks. It covers when to scale vertically and when horizontally, which signals justify a change, how to add nodes safely and in what order, and how to choose between decommission, removenode, replace and assassinate. It also works through the disk arithmetic that decides whether a shrink is safe. The internals of bootstrap streaming, token allocation and zero-copy transfer are covered in Cassandra bootstrap architecture. How tokens map to ranges is covered in token ring and vnodes in depth.

Advertisement

Two axes: node size and node count

There are two independent axes. Vertical scaling (up or down) changes the size of each node: CPU, memory, disk and network. Horizontal scaling (out or in) changes the number of nodes, and with it the fraction of the token ring each node owns. Cassandra prefers horizontal scaling because the partitioner spreads both data and request load evenly across owners. A seven-node data center at RF 3 stores each partition on 3 of 7 machines, so each node holds roughly 3/7 of the data set.

Vertical changes still matter. A node's disk density determines how long it takes to stream, repair and replace it. If you double each node's disk instead of doubling the node count, replacing a failed node takes about twice as long, and a second failure in the same replica set is more likely during that window. Memory is the subtle one. Cassandra's JVM heap should stay modest, often 8 to 31 GB depending on collector and workload (see Java heap tuning). Extra RAM beyond that goes to the OS page cache, which speeds reads only if the hot data set fits. More cores help compaction and request concurrency until the disk becomes the bottleneck.

Signals that justify a change

Change topology when a measured signal says so, not on a calendar. The signals below, and their usual answers, cover most cases:

SignalWhat it meansUsual answer
Disk use per node above about 50% (STCS) or 70% (LCS)Compaction may not have room to rewrite SSTablesScale out, or reduce data (TTL, deletes)
Pending compactions growing for daysWrite volume exceeds compaction throughputMore nodes or faster disks; check compaction strategy
Read p99 rising while CPU idlesPage-cache misses or too many SSTables per readMore RAM per node, or more nodes to shrink each node's data
CPU saturated on all nodes evenlyRequest load exceeds computeScale out, or larger instances
One node hot, others idleHot partition or bad token balanceFix the data model or tokens; adding nodes will not help
Dropped mutations or reads in tpstatsNodes are shedding loadScale out after ruling out GC pauses
Utilisation below 30% for monthsOver-provisionedCandidate for scaling in, after the disk check below

The headroom figures come from compaction. Size-tiered compaction can need as much free space as the SSTables it merges, so a node running STCS at 50% usage can fail a large compaction. Leveled compaction needs far less temporary space but costs more write I/O. Compaction strategies covers the trade in detail. Whatever the strategy, the hot-partition row matters most. Adding nodes to fix a single hot partition only adds cost, because that partition still lives on the same RF replicas.

Advertisement

Scaling out: adding nodes safely

Scale out: add N7Scale in: decommission N6N7 joining (UJ)auto_bootstrapN1..N6 (UN)current ownersstream rangesnodetool cleanup on N1..N6one node at a time, after N7 is UNthenResult: each node owns ~1/7 of ringload per node falls ~14%N6 leaving (UL)still servingN1..N5 (UN)new ownersstream rangesGuard: RF check + disk headroomdecommission refuses below RFbeforeResult: each node owns ~1/5 of ringload per node rises 20%Dead node instead? removenode streams from surviving replicas; replace_address_first_boot reuses its tokensassassinate only edits gossip and moves no data: last resort
Scaling out and in on a six-node, RF 3 data center. Joins stream data to the new node and need cleanup afterwards; a decommission streams data away and raises every survivor's load, so the disk check comes first.

Adding capacity follows a fixed sequence. By default Cassandra refuses to bootstrap a node while another range movement in the same cluster is in progress (the cassandra.consistent.rangemovement property, default true), and you should not override it. Add nodes one at a time:

# 0. Preconditions: every node UN, no streams or repairs running, schema agreed
nodetool status            # all nodes UN
nodetool describecluster   # one schema version
nodetool netstats          # "Not sending any streams"

# 1. New node's cassandra.yaml: same cluster_name, snitch, num_tokens as its peers;
#    seeds point at existing nodes (a new node must not list itself as a seed,
#    or it will skip bootstrap and join empty).
#    On 4.0+, set allocate_tokens_for_local_replication_factor to the DC's RF.

# 2. Start it and watch it join
nodetool status            # new node shows UJ, then UN
nodetool netstats          # bytes streamed per peer

# 3. After it is UN: cleanup on each OLD node, one at a time
nodetool cleanup           # drops ranges the node no longer owns

Do not skip cleanup. Until it runs, the old owners still hold copies of the ranges they gave away. Those copies take disk space and come back to life if the topology later reverts. Cleanup rewrites every SSTable on the node, so it needs disk headroom and competes for I/O. Run it on one node at a time during a quiet period. Bootstrap streaming is throttled by nodetool setstreamthroughput (the yaml key is stream_throughput_outbound from 4.1 and stream_throughput_outbound_megabits_per_sec before that). Raise the limit only if the senders' read latency can absorb the extra load.

Keep num_tokens identical across a data center. Cassandra 4.0 changed the default from 256 to 16 together with the replica-aware token allocator. A new node that silently takes the default from a fresh package can therefore own a very different share of the ring from its 256-token peers. Check nodetool status's Owns column (with a keyspace argument) after every join.

Scaling in: decommission, removenode, replace, assassinate

There are four ways to remove a node, and choosing the wrong one is the most common scaling mistake:

CommandNode stateWho streamsUse when
nodetool decommissionUp, run on the leaving nodeLeaving node sends its ranges to new ownersPlanned shrink of a healthy node
nodetool removenode <host-id>Down, run from any live nodeSurviving replicas send to new ownersNode is dead and will not come back
-Dcassandra.replace_address_first_boot=<ip>Old node down; new node takes its tokensSurviving replicas send to the replacementDead hardware, same capacity wanted
nodetool assassinate <ip>AnyNobody: gossip state onlyA ghost endpoint that removenode cannot clear

Decommission is the safe path for a planned shrink because the leaving node is itself a complete source for its ranges. Since 4.0 it also refuses to proceed if leaving would drop any keyspace below its configured replication factor in that data center, unless you pass --force. Treat --force as a decision to accept fewer replicas than designed. If you really mean to run fewer nodes than RF, first lower the keyspace's RF with ALTER KEYSPACE and run cleanup, then decommission.

Removenode is different. The dead node's data is rebuilt from the other replicas, so any write that reached only the dead node (possible at consistency ONE, or with hints that expired) is lost. Run repair on the affected ranges before removenode if the remaining replicas are reachable, or immediately after. Monitor it with nodetool removenode status. Prefer replace over removenode-then-add when you want the same capacity back, because replace streams the data once instead of twice. Assassinate moves no data at all and only tells gossip to forget an endpoint, so use it last, after removenode has been tried.

Worked example: can six nodes become five?

Take a single data center of six nodes with RF 3, each with a 2 TB data disk. nodetool status reports a load of about 1.2 TB per node, and the main tables use STCS. Utilisation has been low for a quarter, so finance asks whether you can drop to five nodes.

Total stored data is 6 x 1.2 = 7.2 TB, replicas included. Spread over five nodes that becomes 7.2 / 5 = 1.44 TB each, which is 72% of a 2 TB disk. STCS's worst-case compaction can need free space close to the size of the largest tier being merged, and at 72% one large compaction can fill the disk. The answer is not without changes. You would need to switch the large tables to LCS or a capped compaction, shed data, or keep six nodes on smaller disks.

Suppose you switch to LCS and target 70% maximum. Then 1.44 TB is just over the line, and pruning 300 GB of expired data brings it to 6.9 / 5 = 1.38 TB, or 69%. Next, estimate the stream. The leaving node sends about 1.2 TB. At an effective 100 MB/s that is about 3.3 hours, and at a conservative 30 MB/s throttle it is about 11 hours. During that window every request it serves still works, but the cluster has less margin for another failure. Schedule it outside peak hours, confirm the RF guard is satisfied (five nodes is more than RF 3), and record nodetool status before and after. Finally, check that every surviving node's load rose by roughly 20% as predicted. A node that rose much more points at token imbalance worth fixing.

# headroom check before any decommission (run against nodetool status output)
total_load = sum(load_tb for node in dc)          # 7.2
survivors  = len(dc) - 1                          # 5
per_node   = total_load / survivors               # 1.44
limit      = disk_tb * (0.50 if stcs else 0.70)   # 1.0 or 1.4
assert survivors >= rf, "lower RF first"
assert per_node <= limit, f"need {per_node - limit:.2f} TB less data per node"

Multiple data centers, seeds and racks

In a multi-data-center cluster every rule applies per DC. Replication is set per DC, the RF guard counts nodes per DC, and token allocation balances within a DC. Scale one DC at a time, and remember that LOCAL_QUORUM traffic in that DC feels its streaming load while other DCs do not. To retire a whole data center, the order is: stop clients from routing to it, set its RF to 0 with ALTER KEYSPACE (including the system_auth and system_distributed keyspaces), then decommission its nodes. Getting the order wrong makes decommission try to stream the DC's data to its own remaining nodes. Multi-DC topology covers routing and replication placement.

Seeds deserve their own rule. Do not decommission a seed until you have removed it from every node's seed list through a rolling config change, and keep at least two seeds per DC. Rack-aware snitches add one more constraint: keep racks balanced as you add or remove nodes, or one rack ends up owning more replicas than the others.

Failure modes

  • Bootstrap fails halfway. Current versions let you resume with nodetool bootstrap resume. Otherwise wipe the joining node's data directories and start again. Do not leave a half-joined node running.
  • Disk fills during cleanup or after a shrink. Writes fail on that node and compaction stalls. Stop non-essential compactions, add disk, or bring a node back. This is why the arithmetic comes first.
  • Two joins at once. Overriding consistent range movement can leave ranges with inconsistent ownership. Add nodes serially, even when automation makes parallel tempting.
  • Removenode on a node that was only partitioned. If it comes back, it holds tokens that the ring has reassigned. Make sure a dead node stays dead (stop the service, remove it from the load balancer) before removing it.
  • Mixed num_tokens. Ownership skews and the hottest node limits the cluster. Standardise before scaling.
  • Autoscaling from metrics. Elastic autoscaling designed for stateless services fights Cassandra. Each change streams terabytes, so a CPU-driven policy can start a shrink during a traffic dip that ends before the stream does. Automate the runbook, but keep a human or a slow, capacity-based policy in the decision.

What to do next

  1. Record per-node load, disk size, compaction strategy and num_tokens for each DC, and compute current utilisation against the 50% (STCS) or 70% (LCS) line.
  2. Graph pending compactions, read p99, dropped messages and CPU per node, and decide which of the signals above, if any, is firing.
  3. Before scaling in, run the headroom check for N-1 nodes and confirm survivors are at least RF in every DC.
  4. Write the add-node runbook (preconditions, yaml diff, join, verify Owns, cleanup one node at a time) and rehearse it on staging.
  5. Write the remove-node decision tree: healthy means decommission; dead and to be replaced means replace_address_first_boot; dead and not replaced means repair, then removenode; ghost means assassinate.
  6. Set streaming throughput deliberately and estimate stream time from load divided by throughput before every change.
  7. Read bootstrap architecture for what the stream does internally, then repair so removenode never loses acknowledged writes.
Key takeaway: Scale Cassandra horizontally by default, one node at a time, and always run cleanup on the old nodes after a join. Scale in only after the disk arithmetic shows survivors stay under compaction headroom and the node count stays at or above RF in every data center. Use decommission for healthy nodes, replace for dead ones you want back, removenode with repair for dead ones you do not, and assassinate only for ghosts.