Three-way replication is the reason HDFS shrugs off failed disks and dead DataNodes, and it is also the reason teams believe they have a backup when they do not. Replication copies every mistake within seconds. A job that overwrites a partition with empty output, an operator who runs hdfs dfs -rm -r -skipTrash on the wrong path, a compromised admin account, a corrupted NameNode metadata directory or a lost data centre all defeat it, because all three replicas live in the same namespace under the same credentials.

This article is the plan rather than the mechanics. The mechanics of HDFS snapshots and advanced DistCp have their own pages. Here we start from the threats, set recovery objectives per dataset, build three layers of protection (snapshots, incremental replication to a second site, and off-cluster metadata backups), keep data and metadata consistent with each other, and finish with the restore drills that prove any of it works.

What replication does not protect against

ThreatReplicationSnapshotsDistCp to DRMetadata backup
Disk or DataNode failureYesNo needNo needNo need
Accidental delete or bad overwriteNoYes, fastYes, slowerNo
Silent corruption by a buggy jobNoYes, if caught before expiryYes, if lag is longer than detection timeNo
NameNode metadata lossNoNo, snapshots live in that metadataYes, full rebuildYes
Compromised superuser or ransomwareNoNo, superuser can delete themOnly if DR uses separate credentialsOnly if stored offsite
Site or cluster lossNoNoYesYes
Lost KMS keys for encryption zonesNoNoNo, copies are unreadable tooOnly if keys are backed up

Two rows deserve emphasis. Snapshots are metadata inside the NameNode that pin blocks on the same DataNodes, so anything that destroys the namespace or the cluster destroys them too. And for encrypted data, the key server is part of the dataset: a perfect copy of an encryption zone is useless without the keys that decrypt it.

Three layers of protection

Three layers of protection for an HDFS estatePrimary clusterHDFS data+ .snapshot dirsNameNode + JNsfsimage + editsHive metastore DBtables, partitionsKMS, Ranger, KDCkeys, policies, principalsLayer 1: snapshots + trash (same disks)DR HDFS clusterLayer 2: DistCp -update -diff s1 s2changed files onlyOff-cluster metadata vaultLayer 3: fsimage, DB dumps, key backups, configsfetchImagemysqldump / pg_dumpexportReplication protects against disk loss. Only copies on other machines, under other credentials, protect against people.
Layer 1 recovers from mistakes in minutes, layer 2 survives the loss of the cluster, and layer 3 carries the metadata without which copied files are just bytes.

Each layer covers what the one before cannot. Snapshots and trash recover individual mistakes in minutes without moving data. Incremental replication to a DR cluster or object store survives the loss of the primary. The metadata vault holds what the DR copy needs to be usable: Hive table definitions, authorization policies, encryption keys, Kerberos principals and configuration.

Tier your data by RPO and RTO

Not every directory deserves the same treatment, and backing up everything at the strictest tier is how DR budgets get cancelled. Classify datasets by two numbers: the recovery point objective (RPO, how much recent change you can lose) and the recovery time objective (RTO, how long until it is usable again).

TierExamplesRPO / RTOProtection
IrreplaceableRaw event landing zone, regulatory archives1 h / 4 hHourly snapshots, hourly DistCp diff, offsite metadata
Expensive to rebuildCurated warehouse tables24 h / 24 hDaily snapshots and DistCp, metastore dump
DerivableAggregates, feature tables, ML training setsRebuild / daysSnapshots only; keep the pipeline code and inputs safe
Scratch/tmp, staging, user sandboxesNoneTrash only

Write the tier into a manifest that the backup jobs read, so a new dataset starts at a deliberate tier rather than whatever the last engineer copied.

Layer 1: snapshots and trash

Snapshots are read-only, point-in-time views of a snapshottable directory. Creating one copies no blocks; the NameNode records the state, and blocks deleted later stay on disk while any snapshot references them. Pair them with trash, set through fs.trash.interval in minutes, which catches shell deletes but not deletes made with -skipTrash or directly through the FileSystem API. See HDFS trash for those edge cases.

hdfs dfsadmin -allowSnapshot /data/raw/events
hdfs dfs -createSnapshot /data/raw/events h-2026100814
hdfs snapshotDiff /data/raw/events h-2026100813 h-2026100814   # what changed in the hour
hdfs dfs -ls /data/raw/events/.snapshot/                       # list snapshots

Keep a rolling window, such as 48 hourly plus 14 daily snapshots, and expire the rest in the same scheduled job. Watch the cost: a snapshot of a directory that is rewritten daily pins a full extra copy for every day it is retained. Monitor the space each snapshottable root holds that the live tree does not, and alert when it exceeds a budget.

Layer 2: incremental replication with snapshot diffs

A full distcp -update compares every file on both sides, so for a directory with tens of millions of files the listing alone can take hours even when nothing changed. Snapshot-diff mode avoids that: DistCp asks the NameNode for the diff between two snapshots and copies, renames and deletes only what changed. Its preconditions are strict. Both sides must be HDFS, it requires -update, both snapshots must exist on the source, the target must hold a snapshot with the older name, and the target must be unchanged since that snapshot was taken. The cycle that keeps those conditions true:

#!/usr/bin/env bash
set -euo pipefail
SRC=hdfs://prod-nn/data/raw/events
DST=hdfs://dr-nn/data/raw/events
PREV=$(cat /var/lib/dr/events.last)          # e.g. h-2026100813
NEW=h-$(date -u +%Y%m%d%H)

hdfs dfs -createSnapshot "$SRC" "$NEW"
if hadoop distcp -update -diff "$PREV" "$NEW" -pugpbx \
     -m 40 -bandwidth 50 "$SRC" "$DST"; then
  hdfs dfs -createSnapshot "$DST" "$NEW"     # target now equals source at NEW
  echo "$NEW" > /var/lib/dr/events.last
  # Retention is a separate job that expires old snapshots on both sides; it
  # must never delete the name stored in events.last, which the next -diff needs.
else
  # diff precondition broken (target modified, snapshot missing): full resync,
  # copying from the frozen snapshot path so the source cannot move underneath.
  hadoop distcp -update -delete -pugpbx -m 80 "$SRC/.snapshot/$NEW" "$DST"
  hdfs dfs -createSnapshot "$DST" "$NEW"
  echo "$NEW" > /var/lib/dr/events.last
  exit 2                                     # page someone: why did the diff fail?
fi

Three details matter. -pugpbx preserves user, group, permissions, block size and extended attributes; preserving block size keeps the default checksum comparison meaningful between clusters. -bandwidth is in MB per second per map, so 40 maps at 50 caps the job near 2 GB/s. And the DR side must be read-only to everyone except the replication principal; a single analyst writing to the target breaks the diff precondition and forces a full resync. For encryption zones, copy through /.reserved/raw paths so the encrypted bytes and their metadata move unchanged, which requires the same keys to be available at the DR site.

Layer 3: metadata that must travel with the data

Data without metadata is not a recovery. Back up these on their own schedule, to storage that the Hadoop admin credentials cannot delete:

  • NameNode fsimage. hdfs dfsadmin -fetchImage /backup/nn/ downloads the most recent checkpointed image without stopping anything. It is a namespace snapshot as of the last checkpoint, by default at most an hour or a million transactions old. It does not contain blocks: restoring an old image onto the same DataNodes recovers the directory tree, but files whose blocks were deleted since then come back missing, and blocks written later are unknown to the restored NameNode and get deleted. Treat it as last-resort protection against metadata corruption, alongside the HA and JournalNode design in checkpoints and JournalNodes.
  • Hive metastore database. A consistent dump, such as mysqldump --single-transaction or pg_dump, timed right after the data snapshot it should match.
  • KMS keys. Export and escrow the key provider's keystore under separate custody; losing it makes every encryption zone, including the DR copy, unreadable.
  • Ranger policies, Kerberos database and keytabs, and cluster configuration. Without them the DR cluster either refuses everyone or allows everyone.

Keeping data and metadata consistent

Data and metadata are copied at different moments, so a restore must reconcile them. For Hive external tables, the dump taken right after snapshot NEW may reference a partition that landed between the snapshot and the dump. On the DR side, run MSCK REPAIR TABLE to add partitions that exist on disk, and drop registered partitions whose directories are missing. Hive transactional (ACID) tables are harder: their state is spread across base and delta directories and the transaction tables in the metastore, so copying files while compaction runs can give an inconsistent table. Use Hive's own replication commands (REPL DUMP and REPL LOAD) for those, or quiesce compaction during the copy. Ordering helps: snapshot the data, then dump the metastore, then record both identifiers together in a backup manifest so a restore always pairs the right ones. Apache Hive covers the metastore model in more depth.

Worked example: sizing and a real restore

Take a landing zone of 400 TB that changes by about 2 TB a day, replicated over a 10 Gb/s link of which you can use half. Half of 10 Gb/s is about 625 MB/s, so a day's 2 TB moves in roughly 55 minutes, and an hourly diff of about 85 GB moves in under 3 minutes plus job startup. The achievable RPO is the cycle interval plus copy time, a little over an hour, which meets the irreplaceable tier. The initial full copy is the real problem: 400 TB at 625 MB/s takes about 7.5 days, so plan a seeded copy or a phased migration rather than expecting the first run to finish overnight.

Now the incident. At 14:05 a backfill job overwrites /data/raw/events/dt=2026-10-07 with empty files. Detection comes at 14:40. The hourly snapshot from 14:00 predates the damage, so the restore is local:

hdfs dfs -mv /data/raw/events/dt=2026-10-07 /data/quarantine/events-dt=2026-10-07
hadoop distcp -pugpbx \
  /data/raw/events/.snapshot/h-2026100814/dt=2026-10-07 \
  /data/raw/events/dt=2026-10-07

Moving the damaged directory aside, rather than deleting it, preserves evidence and keeps the restore reversible. DistCp parallelises a large restore where hdfs dfs -cp would copy file by file. Then check row counts against the source system before re-enabling downstream jobs. Note the timing trap: the DR side received the empty files at the 15:00 cycle, so if detection had taken three hours, the DR copy would hold the damage too. That is why snapshots on both sides keep a retention window longer than your worst detection time.

Failure modes

  • Broken diff chain. Someone writes to the DR path, or an old snapshot is deleted out of order. The job falls back to a full resync that takes hours; alert on every fallback.
  • Snapshot bloat. Retained snapshots of frequently rewritten data fill the cluster and look like a mystery capacity leak.
  • Checksum mismatches. Different block sizes or checksum types between clusters fail the post-copy check. Preserve block size, or use -skipcrccheck only when you verify another way.
  • Open files. Files still being written when a snapshot is taken are an edge case; Hadoop 3 has dfs.namenode.snapshot.capture.openfiles to freeze their length at snapshot time. Test with your ingestion pattern.
  • Shared credentials. If the primary's superuser can delete DR snapshots, one stolen credential erases both sites.
  • Untested restores. The most common failure is a backup nobody has restored: missing keytabs, a metastore dump of the wrong version, or KMS keys that were never exported.

Trade-offs

  • Second HDFS cluster or object store. A DR cluster supports snapshot-diff DistCp, keeps HDFS semantics and can take over workloads, but it is a second cluster to run. An object store is cheaper and operationally lighter, but DistCp's -diff and -rdiff options do not support object stores, so each run is a full -update copy from a .snapshot/<name> path whose listing cost grows with file count. Immutability then comes from the store's own versioning or object lock rather than from snapshots, and recovery means copying back or running compute against the store.
  • Snapshot frequency against pinned space. Hourly snapshots shrink the window for local recovery but pin more superseded blocks on churning data.
  • Separate credentials against friction. A DR site that primary admins cannot touch survives a compromised account, at the cost of a second approval path for every restore.

What to do next

  1. Write the threat table for your estate and mark which rows you cannot recover from today.
  2. Classify every top-level dataset into a tier with an RPO and RTO, and store it in a manifest the backup jobs read.
  3. Enable snapshots and a retention job on the irreplaceable and expensive tiers; check fs.trash.interval is non-zero.
  4. Stand up the snapshot-diff DistCp cycle to a DR HDFS cluster under separate credentials, with alerts on fallback and on replication lag.
  5. Schedule fetchImage, metastore dumps, KMS key export, Ranger and KDC backups to storage outside the cluster's trust boundary.
  6. Run a restore drill each quarter: one file from a snapshot, one table from DR, and one full-cluster metadata rebuild, timing each against its RTO. Review NameNode HA so the drill also covers failover.
Key takeaway: HDFS replication protects against hardware, not people or sites. Layer snapshots and trash for fast local recovery, snapshot-diff DistCp to a read-only DR target under separate credentials for site loss, and off-cluster backups of fsimage, the Hive metastore, KMS keys, Ranger and Kerberos for everything the data needs to be usable. Tier datasets by RPO and RTO, keep snapshot retention longer than your detection time, pair data and metadata in a manifest, and prove it all with timed restore drills.