Cassandra is famous for staying up when nodes fail, and that reputation leads teams to under-invest in operating it. A cluster that nobody operates does keep serving, for a while. Then an unrepaired replica resurrects deleted data, a disk fills during compaction, or an engineer restarts two nodes in the same rack and a quorum read starts failing. The database's resilience assumes an operator who runs a small number of routine loops correctly and never runs two risky actions at once.

This article describes that operating model as an architecture: the loops, the interfaces they use, the invariants every change must respect, and the procedures for the node lifecycle, upgrades and backups. Neighbouring pages go deeper on individual loops. Operational metrics covers what to watch, repair covers anti-entropy, and streaming and bootstrap covers how data moves when topology changes. The version-specific details here are for Cassandra 4.x and 5.0.

Advertisement

The operating model: four loops

Every operational task falls into one of four loops, and they differ in their risk and in how often they run.

  • Observe runs continuously: metrics, logs, and ad hoc state queries. It produces the health signal the other loops depend on.
  • Maintain runs on a schedule: repair within every gc_grace_seconds window, cleanup after topology changes, snapshot rotation, and compaction tuning.
  • Change runs on demand: add or remove capacity, restart for config or OS patches, upgrade versions, alter schema.
  • Recover runs on failure: replace a dead node, restore data, rebuild a datacenter.

The architectural rule is that Change and Recover actions are serialised through one orchestrator (a script, a Kubernetes operator or a workflow engine), which takes a cluster-wide lock and refuses to proceed unless the Observe loop reports a healthy cluster. Most multi-node incidents come from two changes overlapping, for example a repair streaming data while a node is decommissioning, or a restart during a major upgrade.

Operations as four control loops around the clusterCassandra clusternodes, racks, DCsObservemetrics, logs, virtual tablesMaintainrepair, cleanup, compactionChangescale, restart, upgrade, configRecoverreplace, restore, rebuildJMX, CQLnodetoolone node at a timestreamingOrchestratorpreconditions, locksEvery change and recovery action is gated by the same health checks the observe loop produces.
The orchestrator is the only path to Change and Recover actions, and each action is gated by the health signal from the Observe loop.

The operator's interfaces

Four interfaces expose the cluster, and a mature setup uses each for what it is good at.

InterfaceScopeUse it for
nodetool (over JMX)One node per callLifecycle actions, flush, drain, snapshots, throughput throttles, status
Virtual tables (system_views, 4.0+)One node per queryReading settings, thread pools, clients and running tasks over CQL
cassandra.yaml and JVM optionsOne node, at restartDurable configuration, managed by config management
Logs (system.log, debug.log)One nodeGC pauses, dropped messages, streaming and compaction errors
-- Cassandra 4.0+ virtual tables: per-node state over CQL, no JMX needed
SELECT name, value FROM system_views.settings WHERE name = 'compaction_throughput';
SELECT name, active_tasks, pending_tasks, blocked_tasks FROM system_views.thread_pools;
SELECT address, username, driver_name, driver_version FROM system_views.clients;
SELECT keyspace_name, table_name, kind, progress, total FROM system_views.sstable_tasks;

Almost every interface is per node. There is no cluster-wide admin API in the core database, so any "cluster" operation is a loop over nodes. Your tooling must record which nodes it touched and in what state it left them. Keep cassandra.yaml in version control and render it per node, because a setting changed with nodetool at runtime, such as setcompactionthroughput, is lost at the next restart.

Advertisement

Preflight: the invariant every change checks

Before any change, and again after each step, confirm four things. Every node is Up/Normal (UN in nodetool status). The cluster agrees on one schema version. There is no large compaction backlog. And no streaming is in progress. Encode this as a script so that it is impossible to skip under pressure:

#!/usr/bin/env bash
# preflight.sh - refuse to start a change unless the cluster is quiet and healthy.
set -euo pipefail
fail() { echo "PREFLIGHT FAIL: $*" >&2; exit 1; }

status=$(nodetool status)
echo "$status" | grep -E '^(DN|UL|UJ|UM|DL|DJ|DM) ' && fail "a node is not Up/Normal"

versions=$(nodetool describecluster | sed -n '/Schema versions:/,/^$/p' | grep -c '\[')
[ "$versions" -eq 1 ] || fail "schema disagreement ($versions versions)"

pending=$(nodetool compactionstats | awk '/pending tasks/ {print $3}')
[ "${pending:-0}" -lt 50 ] || fail "compaction backlog: $pending pending tasks"

nodetool netstats | grep -q 'Not sending any streams' || fail "streaming in progress"
echo "preflight ok"

The thresholds are placeholders. Take them from your own baseline, because a cluster using leveled compaction under steady write load may always show some pending tasks. The point is that the orchestrator runs this, and a human cannot talk it out of a failure.

The node lifecycle

Topology changes are where data moves, so each has a specific command and specific preconditions.

SituationActionWhat happens
Add capacityStart new node with auto_bootstrap on (the default), not listed as a seedIt streams its token ranges from existing replicas, then joins as UN
After addingnodetool cleanup on each old node, one at a timeDeletes data the old node no longer owns; disk is not reclaimed until you do this
Remove a live nodenodetool decommission on that nodeIt streams its ranges to the new owners, then leaves
Remove a dead nodenodetool removenode <host-id> from any live nodeRemaining replicas re-stream the dead node's ranges
Replace a dead nodeStart a fresh node with -Dcassandra.replace_address_first_boot=<dead-ip>It takes over the dead node's tokens and streams their data
Node will not leavenodetool assassinate <ip>Last resort; drops gossip state without streaming, so repair afterwards

Worked example. A six-node cluster in one datacenter, three racks, replication factor 3, and each node holds about 1 TB. Node 4 loses its disk. With RF 3 and one replica per rack, every range still has two live replicas, so LOCAL_QUORUM reads and writes continue. Hints for node 4 accumulate on its peers, but only for max_hint_window (3 hours by default), after which writes to it are dropped from hinting.

The right move is replacement, not removal plus addition, because replacement streams the data once and leaves every other node's token ownership unchanged. Provision a new node in the same rack, set replace_address_first_boot to node 4's address, and start it. It streams about 1 TB from the surviving replicas, throttled by stream_throughput_outbound, whose default in 5.0 is 24 MiB/s per sending node. At that rate, a terabyte takes many hours unless you raise the throttle. Raise it deliberately while watching client latency. Because the node was down longer than the hint window, run a repair of its ranges once it is UN. If a dead node has been gone longer than gc_grace_seconds (10 days by default), never simply bring the old one back, because its stale data can resurrect deletes. Wipe it and replace it.

Rolling restarts and configuration changes

A restart is the most frequent change, so make it boring. Drain first so memtables flush and the commit log does not need replaying. Restart one node at a time, and never start the next until the previous one is UN, serving clients and has had time to receive hints. With racks mapped to availability zones, restarting nodes rack by rack is safe for LOCAL_QUORUM at RF 3, but restarting two nodes in different racks at once is not.

# Rolling restart of one node (run from the orchestrator, one node per step)
./preflight.sh
nodetool disablebinary           # stop taking client connections
nodetool drain                   # flush memtables, stop accepting writes
sudo systemctl restart cassandra
until nodetool status | grep -q "^UN  $NODE_IP "; do sleep 10; done
until nodetool info | grep -q 'Native Transport active: true'; do sleep 5; done
sleep 60                         # let peers deliver hints and caches warm
./preflight.sh                   # the next node starts only if this passes

Configuration changes follow the same procedure, with one node as a canary. Change it, restart, watch latency and GC for a meaningful period, then roll through the rest.

Upgrades

Patch upgrades within a release line are rolling restarts with a new binary. Major upgrades need a plan, because for the duration of the roll the cluster runs mixed versions and some operations are unsafe. The widely followed rules during a mixed-version window are: no topology changes, no repairs, and no schema changes. Keep the window short.

The sequence is to read the release's NEWS.txt for removed settings and changed defaults, and test the upgrade on a restored copy of production data. Snapshot every node, then upgrade one node, check it, and roll through the rest. Finally, run nodetool upgradesstables where the release notes require SSTables to be rewritten into the new format. Configuration names changed in 4.1 (for example compaction_throughput_mb_per_sec became compaction_throughput with explicit units). Older names are still accepted for compatibility, but render the new form from config management so that the next upgrade does not surprise you. Drivers need attention too: confirm the protocol versions your drivers negotiate against the new server before the roll, not after.

Backups and restore

Replication is not backup. It faithfully replicates a bad DROP or TRUNCATE, or a buggy batch job, to every replica. Cassandra's native backup is the snapshot: nodetool snapshot flushes memtables and hard-links every live SSTable into a snapshots/<tag> directory. It is near-instant and initially costs no space, but the links pin files that compaction would otherwise delete, so disk use grows until you clear the snapshot. With auto_snapshot enabled (the default), dropping or truncating a table snapshots it first, which is a valuable safety net that also consumes disk.

# Snapshot every node at roughly the same time with a shared tag
TAG=daily-$(date +%Y%m%d)
nodetool snapshot -t "$TAG" shop                     # hard links, near-instant
# upload .../data/shop/<table>-<id>/snapshots/$TAG/ to object storage, then:
nodetool clearsnapshot -t "$TAG" shop                # hard links pin disk space

# Restore one table on the same topology (4.0+): stage files, then import live
nodetool import shop orders /restore/shop/orders

For point-in-time recovery between snapshots, incremental_backups: true hard-links each newly flushed SSTable into a backups/ directory, which you ship and delete in turn. Commit log archiving can narrow the gap further, at operational cost.

Restore is where plans fail, so rehearse it. Restoring to the same topology is straightforward: copy each node's files back and load them with nodetool import (4.0+) or nodetool refresh. Restoring to a different topology needs sstableloader, which streams the rows to whichever nodes now own them. Record each node's tokens alongside its snapshot, because a restore onto different tokens without a loader will serve the wrong data. Time a full restore at least quarterly. That timing is your real recovery time objective, not the one in a design doc. The wider planning is covered in the disaster recovery guide.

Capacity and headroom

Compaction needs free disk to write merged SSTables before it deletes the inputs. Size-tiered compaction can temporarily need as much free space as the SSTables being merged, which on a large table approaches half the disk. Leveled and unified compaction need much less, but snapshots, streaming and repair all consume space too. A common planning rule is to add capacity well before disks reach 50 to 60 percent under size-tiered compaction. Remember that adding a node does not reduce disk use on old nodes until cleanup runs.

Also track data per node. The practical limit is less about disk and more about how long it takes to stream a replacement. A node holding 4 TB takes days to rebuild at conservative throttles, and during that time the cluster runs with reduced redundancy on those ranges. That window is the reason many operators cap density per node.

Automating it safely

Mature teams codify these procedures in an orchestrator rather than in runbook prose. That can be a set of scripts with a cluster lock, a tool dedicated to repair scheduling such as Cassandra Reaper, or a Kubernetes operator such as those from the K8ssandra project that manage restarts and scaling through custom resources. Whatever you choose, keep four properties: one change at a time per cluster, preflight before and after every step, per-rack or per-zone ordering, and an audit log of which node is in which state so that an interrupted run can resume. When automation stops on a failed check, a human should decide the next step. Automation should never retry its way through a failing check.

Failure modes and trade-offs

  • Overlapping changes. A repair and a decommission, or two restarts in different racks, remove the redundancy each assumed. Serialise through one lock.
  • Forgotten cleanup. After scaling out, old nodes keep data they no longer own. Disk looks full and backups grow.
  • Snapshots that are never cleared. Hard links pin deleted SSTables until the disk fills. Alert on snapshot size.
  • Repair outside gc_grace. Deleted data reappears once tombstones are purged on one replica but not the others.
  • Runtime-only tuning. A nodetool throttle that fixed an incident vanishes at the next restart and the incident returns.
  • Density versus recovery time. Fewer, larger nodes are cheaper to run and slower to replace. Choose density from the rebuild time you can tolerate, not the disk you can buy.
  • Automation versus judgement. Orchestrators remove fat-finger errors but can apply a bad change quickly. A canary node and a hard stop on failed checks keep the speed without the blast radius.

What to do next

  1. Write the preflight script for your cluster, with thresholds from a week of baseline data, and make every change run it.
  2. Put cassandra.yaml and JVM options under config management, and remove every runtime-only tweak.
  3. Schedule repair so that every table completes within gc_grace_seconds, with margin.
  4. Rehearse a node replacement with replace_address_first_boot and record how long streaming takes.
  5. Automate daily snapshots with shared tags, off-node upload and clearsnapshot, and time a full restore each quarter.
  6. Document the mixed-version rules for your next major upgrade before you start it.
Key takeaway: Operating Cassandra well means running four loops: observe, maintain, change and recover. Every change and recovery action goes through one orchestrator that checks node state, schema agreement, compaction backlog and streaming before and after each step. Replace dead nodes rather than removing and re-adding them. Run cleanup after scaling out, repair within gc_grace_seconds, and roll restarts and upgrades one node and one rack at a time. Treat snapshots as backups only once you have shipped them off the node and timed a restore. Choose node density from the rebuild time you can tolerate.