Teams test backups by checking that the upload job succeeded. That proves the bytes left the node. It does not prove that you can put a cluster back in the time your business expects, onto the hardware you will actually have, with data you have checked. This article is about the restore half: what a backup must carry to be restorable, how to estimate restore time before an incident, how to restore onto new nodes or a different node count, how to recover a few bad rows without rolling back the whole table, and how to verify and drill all of it.

The mechanics of snapshots, incremental backups, commit log archiving and the basic restore commands are covered in Cassandra Backups, in depth. Read that first if nodetool snapshot and nodetool import are new to you; this page builds on them. Commands assume Cassandra 4.0 or later.

Advertisement

Start from the restore, not the backup

Design backups backwards from restores. List the incidents you expect, then decide what each restore must look like. The method changes the time it takes far more than the backup format does.

IncidentRestore shapeWhat dominates the time
One table truncated or droppedLocal snapshot, nodetool import on every nodeNumber of nodes you can work on at once
Bad application writes since 09:00Old snapshot into a side cluster, copy corrected rows backFinding the affected partitions
Whole cluster lost, same shape rebuiltDownload per node, pinned tokens, importDownload bandwidth per node
Cluster lost, rebuilt smaller or largerDownload all, sstableloader into new clusterStreaming and compaction of RF copies

Two rules fall out of the table. First, restores of a single table should never need the object store, because Cassandra already takes an automatic snapshot before TRUNCATE and DROP while auto_snapshot is on. Second, whole-cluster restores are bandwidth problems, so the number that matters is gigabytes per second into the new nodes, not the size of the backup.

What a restorable backup must contain

A snapshot directory is just SSTables plus a schema.cql and a manifest.json. That is not enough to rebuild a cluster. To restore onto fresh nodes with the same layout you also need every node's tokens, its data centre and rack, the Cassandra version that wrote the files, and the exact moment the snapshot started on each node. Capture them at backup time, in the same object-store prefix as the files, because during an incident the old cluster is gone and cannot tell you.

#!/usr/bin/env bash
# Run on each node right after `nodetool snapshot -t "$TAG"`; uploads beside the SSTables.
set -euo pipefail
TAG="$1"; HOST=$(hostname -f); OUT=/var/tmp/restore-manifest-$TAG.json
TOKENS=$(cqlsh -e "SELECT tokens FROM system.local;" | sed -n '4p' | tr -d " {}'")
DC_RACK=$(cqlsh -e "SELECT data_center, rack FROM system.local;" | sed -n '4p' | tr -d ' ')
VERSION=$(nodetool version | awk '{print $2}')
cat > "$OUT" <<EOF
{"host": "$HOST", "tag": "$TAG", "taken_at": "$(date -u +%FT%TZ)",
 "cassandra_version": "$VERSION", "dc_rack": "$DC_RACK", "tokens": "$TOKENS"}
EOF
cqlsh -e "DESCRIBE SCHEMA" > /var/tmp/schema-$TAG.cql
aws s3 cp "$OUT" "s3://backups/cass-prod/$TAG/$HOST/restore-manifest.json"
aws s3 cp /var/tmp/schema-$TAG.cql "s3://backups/cass-prod/$TAG/$HOST/schema.cql"

The token list is the critical field. With 16 virtual nodes per host it is 16 numbers; with the old default of 256 it is 256. Without them, a same-shape restore turns into an sstableloader restore, which is several times slower. How tokens map ranges to nodes is explained in the token ring and vnodes.

Advertisement

RTO arithmetic: a worked example

Recovery time objective is a promise, so compute it rather than hoping. Take a production cluster with 12 nodes, replication factor 3, and 1.5 TB of live data per node after compression, which is 18 TB on disk and about 6 TB of unique data.

Object storesnapshots + manifestStaging disksdownload + verifyRestore methodimport / sstableloader / side clusterfetchloadSame tokensnodetool import per nodeNew topologysstableloader streamsSide clustercopy chosen rows backVerifyrow samples, digests, repair, app smoke testRecord timingsfeed the next RTO estimate
Every restore is fetch, load, verify, record. The load step is the only part that changes with the scenario; the timings you record are what make the next RTO estimate honest.

Same-shape restore. Each new node downloads only its own 1.5 TB. If the object store gives each node 400 MB/s, that is about 1.05 hours of download per node, and all 12 run in parallel. nodetool import -cd then copies the staged files into the table and verifies them at local-disk speed, so budget tens of minutes per node, not hours. With schema creation, node start-up and a smoke test, an honest RTO is around 2 to 3 hours.

Different-shape restore. Now the new cluster has 8 nodes, so tokens differ and you must stream. The loader reads every node's files, and each row exists in 3 nodes' files. It sends every copy it reads to all 3 current replicas, so the cluster receives about 18 TB times 3, or 54 TB of writes, which compaction later reduces back to 18 TB. At an aggregate 2 GB/s of ingest across the new cluster, 54 TB is 7.5 hours of streaming before compaction catches up. Loading only one replica's worth of files would cut this to a third, but working out a non-overlapping set of files is error prone with vnodes, so most teams accept the amplification and plan for it.

Logical restore of bad rows. The time is mostly human: finding the keys. That is why the incident runbook for application bugs should start with "which partitions?" rather than "which backup?".

Write these numbers down per keyspace. If the product owner was promised one hour, you have learned that the promise needs either smaller nodes, faster storage or a standby data centre.

Restoring onto new nodes with the same shape

When you rebuild with the same number of nodes, pin each new node to an old node's tokens before its first start. In cassandra.yaml on new node 4, set initial_token to the comma-separated list from old node 4's manifest and make num_tokens equal the number of tokens in that list. Match the data centre and rack too, or replica placement under NetworkTopologyStrategy will differ even with identical tokens.

# 1. On every new node, before first start (values from that twin's restore-manifest.json)
num_tokens: 16
initial_token: -9112334456123,-7741852002931,...   # 16 values, exactly as saved

# 2. Start seeds, then the rest; the cluster is empty. Apply the saved schema once:
cqlsh -f schema.cql

# 3. On each node, download its twin's snapshot into a staging dir laid out as ks/table/
aws s3 sync s3://backups/cass-prod/$TAG/$TWIN/app_ks/orders/ /restore/app_ks/orders/

# 4. Import, keeping the staging copy until verification passes
nodetool import -cd app_ks orders /restore/app_ks/orders

Two details matter. The restored table has a new table ID, so its directory name differs from the one in the backup; that is fine because import takes the source directory as an argument. And system_auth roles live in a table like any other, so either restore them the same way or recreate roles from your configuration management before applications reconnect.

Restoring into a different topology

When the shape changes, use the bulk loader. It acts as a client: it reads each SSTable, works out the current owner of every partition from the new cluster's ring, and streams the data there. The last two path components must be the keyspace and table names.

# Run from a loader host (or several in parallel, each with a different subset of old nodes' files)
for OLD in node01 node02 node03 node04; do
  aws s3 sync "s3://backups/cass-prod/$TAG/$OLD/app_ks/orders/" "/load/$OLD/app_ks/orders/"
  sstableloader -d new-seed-1,new-seed-2 "/load/$OLD/app_ks/orders"
done

Run several loader hosts so the loader is not the bottleneck, and use the loader's throttle option to keep the new cluster responsive; check sstableloader --help on your version, because the units of the throttle flags have changed across releases. Expect pending compactions to spike during and after the load. Check nodetool compactionstats and do not call the restore finished until pending tasks are near zero, because read latency stays poor until then. Streaming behaviour and its tuning knobs are covered in bootstrap and streaming.

Row-level recovery from a side cluster

The most common real restore is not a disaster. A deploy at 09:00 wrote wrong prices into a few thousand rows of app_ks.orders. Rolling the whole table back would lose every correct order since the snapshot. Instead, restore last night's snapshot into a small side cluster, read the old versions of the affected rows, and write them back into production.

The trick is timestamps. Cassandra resolves conflicts by write timestamp, so if you write the old values with their old timestamps, the bad writes from 09:00 win again. Write the corrections with a timestamp just after the bad write, but only if the current value was written inside the one-hour bad window. The script below does exactly that, using the WRITETIME function to compare.

# Copy corrected rows from a side cluster back to production, without clobbering newer writes.
from cassandra.cluster import Cluster
side = Cluster(["side-1"]).connect("app_ks")
prod = Cluster(["prod-1"]).connect("app_ks")
BAD_START_US = 1791018000000000          # 09:00 UTC on the incident day, microseconds
get_old = side.prepare("SELECT price FROM orders WHERE order_id = ?")
get_cur = prod.prepare("SELECT price, WRITETIME(price) AS wt FROM orders WHERE order_id = ?")
fix = prod.prepare("UPDATE orders USING TIMESTAMP ? SET price = ? WHERE order_id = ?")

for key in open("affected_keys.txt"):
    oid = key.strip()
    old = side.execute(get_old, [oid]).one()
    cur = prod.execute(get_cur, [oid]).one()
    if old is None or cur is None:
        continue                            # row created after the snapshot: handle by hand
    if cur.wt < BAD_START_US:
        continue                            # never touched by the bad deploy
    if cur.wt > BAD_START_US + 3_600_000_000:
        continue                            # changed again after the bad window: leave it
    prod.execute(fix, [cur.wt + 1, old.price, oid])

Use a side cluster rather than importing old files into production. An old snapshot may predate tombstones that compaction has since purged, and importing it would bring deleted rows back. The side cluster contains that risk.

Verifying a restore

A restore is finished when you have evidence that the data is right, not when the import command returns. Use three layers of checks, from cheap to expensive.

  1. Structural. Every expected table exists, nodetool tablestats shows non-zero SSTables on every node, and the estimated partition counts are within a few percent of the figures recorded at backup time. Estimates are approximate, so treat a 2 percent gap as normal and a 30 percent gap as a missing node's files.
  2. Sampled digests. Pick a few thousand random partition keys from the backup's own data (for example from the manifest of known keys, or with SELECT ... WHERE token(pk) > ? LIMIT on random tokens in the source), read them at QUORUM from the restored cluster, and compare a hash of each row against the same hash computed on the side copy. Any mismatch fails the drill.
  3. Application. Point a staging build of the application at the restored cluster and run its read-only smoke tests. This catches schema and role problems that row checks miss.

Then run nodetool repair -pr on each node, because per-node snapshots were not taken at exactly the same instant. Repair concepts are in Cassandra repair.

Restore drills as code

Restores fail in drills for boring reasons: an expired credential, a manifest field nobody wrote, a version mismatch. Find those monthly, automatically. If you use Medusa, its commands map onto the same steps; its documented subcommands include backup-cluster, list-backups, verify --backup-name NAME (add --enable-md5-checks for content checks), restore-cluster, restore-node and report-last-backup. Read its restore documentation for the host-mapping options before relying on a cross-topology restore.

#!/usr/bin/env bash
# Monthly restore drill: rebuild a scratch cluster from last night's backup and assert timings.
set -euo pipefail
START=$(date +%s)
TAG=$(aws s3 ls s3://backups/cass-prod/ | awk '{print $2}' | sort | tail -1 | tr -d /)
./provision-scratch-cluster.sh --nodes 12 --tokens-from "s3://backups/cass-prod/$TAG"
./restore-all-nodes.sh "$TAG"                   # download + import, per node in parallel
LOADED=$(date +%s)
python3 verify_sample.py --keys 5000 --source side --target scratch   # exits 1 on mismatch
END=$(date +%s)
echo "tag=$TAG load_s=$((LOADED-START)) total_s=$((END-START))" | tee -a drill-history.log
[ $((END-START)) -lt 10800 ] || { echo "RTO budget of 3h exceeded"; exit 1; }
./destroy-scratch-cluster.sh

Graph the drill history. Restore time grows with data, so a slowly rising drill time warns you that the RTO promise will break.

Failure modes

FailureSymptomPrevention
Missing tokensForced into a slow sstableloader restoreWrite the manifest with every snapshot; fail the backup if it is missing
Wrong rack or DC on new nodesReplicas land on the wrong nodes, reads miss dataRestore cassandra-rackdc.properties from the manifest
Version mismatchImport refuses files or nodes fail to startRecord the version; restore onto the same major version, then upgrade
Import moved the only copyA failed import leaves nothing to retryAlways -cd, keep staging files until verified
Compaction backlog after loadingHigh read latency for hoursBudget it in RTO; raise compaction throughput during restore
Resurrected deletesDeleted rows reappearOld backups only into side clusters
Credentials deleted with the clusterBackups unreadableSeparate account, documented break-glass access

Trade-offs

Same-shape restores are fast but tie you to the old node count and token layout; that is acceptable for disaster recovery and awkward for capacity changes, where restoring then expanding or shrinking with normal operations is usually faster than loading into the new shape. Bulk loading is flexible but writes about RF times the data and leaves a compaction debt. Logical restores are precise but need people who can identify affected keys, so invest in audit columns or change logs that make that question answerable. Standby data centres give the shortest RTO of all, at the cost of a second copy of the cluster, and they do not protect against bad writes, which replicate to them instantly.

What to do next

  1. Write the incident table above for your own keyspaces and assign an RPO and RTO to each.
  2. Add the restore manifest script to your snapshot job and fail the job if the manifest is missing.
  3. Compute RTO for a same-shape and a different-shape restore using your real node sizes and measured download speed.
  4. Write and rehearse the logical restore script against a staging table with a simulated bad deploy.
  5. Build the monthly drill, include sampled digest verification, and keep its timing history.
  6. Review SSTable internals so you know what you are moving around during a restore.
Key takeaway: A Cassandra backup is only as good as the restore you have timed. Capture tokens, rack, version and schema with every snapshot, compute restore time from bandwidth and replication factor, pin tokens for same-shape rebuilds, use the bulk loader when the shape changes, recover bad rows through a side cluster with write timestamps, verify with sampled digests and repair, and run the whole thing as a monthly drill.