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.
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.
| Incident | Restore shape | What dominates the time |
|---|---|---|
| One table truncated or dropped | Local snapshot, nodetool import on every node | Number of nodes you can work on at once |
| Bad application writes since 09:00 | Old snapshot into a side cluster, copy corrected rows back | Finding the affected partitions |
| Whole cluster lost, same shape rebuilt | Download per node, pinned tokens, import | Download bandwidth per node |
| Cluster lost, rebuilt smaller or larger | Download all, sstableloader into new cluster | Streaming 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.
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.
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/ordersTwo 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"
doneRun 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.
- Structural. Every expected table exists,
nodetool tablestatsshows 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. - 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) > ? LIMITon random tokens in the source), read them atQUORUMfrom the restored cluster, and compare a hash of each row against the same hash computed on the side copy. Any mismatch fails the drill. - 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.shGraph the drill history. Restore time grows with data, so a slowly rising drill time warns you that the RTO promise will break.
Failure modes
| Failure | Symptom | Prevention |
|---|---|---|
| Missing tokens | Forced into a slow sstableloader restore | Write the manifest with every snapshot; fail the backup if it is missing |
| Wrong rack or DC on new nodes | Replicas land on the wrong nodes, reads miss data | Restore cassandra-rackdc.properties from the manifest |
| Version mismatch | Import refuses files or nodes fail to start | Record the version; restore onto the same major version, then upgrade |
| Import moved the only copy | A failed import leaves nothing to retry | Always -cd, keep staging files until verified |
| Compaction backlog after loading | High read latency for hours | Budget it in RTO; raise compaction throughput during restore |
| Resurrected deletes | Deleted rows reappear | Old backups only into side clusters |
| Credentials deleted with the cluster | Backups unreadable | Separate 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
- Write the incident table above for your own keyspaces and assign an RPO and RTO to each.
- Add the restore manifest script to your snapshot job and fail the job if the manifest is missing.
- Compute RTO for a same-shape and a different-shape restore using your real node sizes and measured download speed.
- Write and rehearse the logical restore script against a staging table with a simulated bad deploy.
- Build the monthly drill, include sampled digest verification, and keep its timing history.
- Review SSTable internals so you know what you are moving around during a restore.