Deciding to leave Hadoop is the easy part. The article on why Hadoop is declining covers the reasons and the strategic mistakes. This article is about execution: how to move petabytes of files, thousands of table definitions, hundreds of scheduled jobs and a security model from a cluster that is still in production to cloud object storage, without losing data and with a way back if something goes wrong.
The central idea is to treat the migration as four separate lanes, data, metadata, jobs and security, which progress independently and meet at a per-table cutover gate. Teams that move everything in one weekend find out on Monday which of the four they forgot.
Choose the target shape first
Three target shapes exist, and the choice changes every later step:
| Shape | What moves | Good for | Cost |
|---|---|---|---|
| Lift and shift | HDFS on cloud VMs, same Hive and Spark versions | Hard deadlines such as a data-centre exit | Keeps the coupled storage and compute you were leaving |
| Managed Hadoop on object storage | Data to S3, GCS or ADLS; jobs on EMR, Dataproc or HDInsight | Large Spark and Hive estates with few rewrites | Path rewrites, committer and listing behaviour change |
| Re-platform to a lakehouse | Data to object storage as Iceberg tables; Spark and Trino; a new orchestrator | Long-term cost and flexibility | Most work up front: table conversion and job porting |
Most organisations combine the last two: managed Spark on object storage first, converting the most-used tables to Iceberg as they move. Pure lift and shift is rarely worth it, because cloud VM disks priced for HDFS triple replication are expensive and you still run NameNodes.
Inventory from the fsimage, not from du
Running hdfs dfs -du across a large namespace loads the NameNode and still tells you nothing about file counts or age. The fsimage, the NameNode's checkpoint of the namespace, has everything, and the Offline Image Viewer turns it into a table you can analyse without touching the live cluster:
# 1. Pull the latest checkpointed image and convert it to TSV (offline; no NameNode load)
hdfs dfsadmin -fetchImage /tmp/fsimage/
hdfs oiv -p Delimited -delimiter $'\t' -i /tmp/fsimage/fsimage_000000000123456789 -o /tmp/fsimage.tsvA short script then sizes each dataset by logical bytes, file count, small files and idle time:
# 2. Size every top-level dataset: logical bytes, files, small files, last modification
import csv, collections, datetime as dt
stats = collections.defaultdict(lambda: {"bytes": 0, "files": 0, "small": 0, "newest": 0})
with open("/tmp/fsimage.tsv") as f:
rows = csv.reader(f, delimiter="\t")
header = next(rows)
i = {name: n for n, name in enumerate(header)}
for r in rows:
if r[i["Permission"]].startswith("d"):
continue # directories carry no bytes
parts = r[i["Path"]].split("/")
key = "/".join(parts[:4]) # e.g. /warehouse/sales.db/orders
s = stats[key]
size = int(r[i["FileSize"]])
s["bytes"] += size # logical size, before replication
s["files"] += 1
s["small"] += size < 16 * 1024 * 1024
mtime = dt.datetime.strptime(r[i["ModificationTime"]], "%Y-%m-%d %H:%M")
s["newest"] = max(s["newest"], mtime.timestamp())
for key, s in sorted(stats.items(), key=lambda kv: -kv[1]["bytes"]):
age_days = (dt.datetime.now().timestamp() - s["newest"]) / 86400
print(f'{key}\t{s["bytes"]/1e12:.2f} TB\t{s["files"]}\t{s["small"]}\t{age_days:.0f}d idle')The inventory drives three decisions. Datasets idle for a year are candidates for archive storage or deletion rather than migration. Directories with millions of small files should be compacted before the copy, because object-store requests are charged and slow per object; see the small files article. The largest active tables set the critical path.
Transfer arithmetic
Size the copy in logical bytes. HDFS stores three replicas by default, so a cluster that reports 2 PB of raw usage holds roughly 670 TB of data; erasure-coded directories differ, and the fsimage FileSize column is already logical. Object stores handle their own durability, so you copy each byte once.
Then divide by sustainable bandwidth. A 10 Gbit/s link carries 1.25 GB/s at best; plan for about 80 percent, around 1 GB/s or 86 TB per day. The 670 TB estate therefore needs about eight days of continuous copying, assuming the cluster can read that fast while serving production. In practice the copy runs in a capped YARN queue during business hours and faster at night, so plan for two to three weeks and add incremental passes. If the arithmetic gives months, consider offline transfer appliances for cold data and keep the network for active tables.
The data lane: bulk then incremental
DistCp is the workhorse; its internals and the expensive object-store commit phase are explained in the advanced DistCp article. For a migration, three details matter.
# Bulk copy one partitioned table, skipping CRC comparison across filesystems
hadoop distcp \
-Dmapreduce.job.queuename=migration \
-update -skipcrccheck -numListstatusThreads 40 \
-m 200 -bandwidth 50 \
hdfs://prod-nn/warehouse/sales.db/orders \
s3a://acme-lake/warehouse/sales.db/orders
# Incremental: copy only partitions written since the last pass
for d in $(hdfs dfs -ls /warehouse/sales.db/orders | awk '$6 >= "2026-09-28" {print $8}'); do
hadoop distcp -update -skipcrccheck "hdfs://prod-nn$d" "s3a://acme-lake${d}"
done-skipcrccheckis needed because HDFS block checksums are not comparable with object-store checksums, so DistCp falls back to size comparison. Validate content separately, as shown below.-diffwith HDFS snapshots is the efficient way to copy only changes, but it requires the target to be a snapshottable HDFS directory holding the same base snapshot. An object store cannot be that target, so use-updateor copy by partition.- Most warehouse tables are append-by-partition. Copying only new date partitions on each pass is cheaper and more predictable than repeated full listings, and it gives a clear catch-up boundary for the cutover.
The metadata lane: catalog and table format
Files without a catalog are not tables. Export the Hive Metastore schema and register tables in the target catalog, whether that is AWS Glue, Dataproc Metastore, a Hive Metastore on the cloud or an Iceberg REST catalog. The simplest path keeps tables as Hive tables with locations rewritten from hdfs:// to s3a:// or gs://. The better path converts the important tables to Iceberg, using the Spark procedures Iceberg provides:
-- Spark SQL with the Iceberg runtime and a session catalog over the Hive Metastore
-- a) Trial: an Iceberg table that READS the Hive table's files, leaving the source untouched
CALL spark_catalog.system.snapshot('sales.orders', 'sales.orders_ice_trial');
-- b) In place: replace the Hive table with an Iceberg table over the same files
CALL spark_catalog.system.migrate('sales.orders');
-- c) Existing Iceberg table: register files already copied to the object store
CALL lake.system.add_files(
table => 'sales.orders',
source_table => '`parquet`.`s3a://acme-lake/warehouse/sales.db/orders`'
);snapshot creates a trial Iceberg table that reads the existing files without changing the source, so you can test queries before committing. migrate replaces the Hive table in place. add_files registers files into an existing Iceberg table. Iceberg removes the directory-listing dependence that makes Hive tables slow on object storage, and it gives atomic commits; the Iceberg and Trino article covers operating those tables afterwards.
The jobs lane: paths, committers and orchestration
Job porting is usually the longest lane. Hard-coded HDFS URIs must become configuration. Jobs that rename output directories to publish results need attention, because a rename on S3 is a copy plus delete per object, not an atomic metadata operation. For Spark writing plain files to S3, use the S3A committers through the spark-hadoop-cloud module, setting spark.hadoop.fs.s3a.committer.name to magic or directory and the commit protocol to org.apache.spark.internal.io.cloud.PathOutputCommitProtocol, with the matching binding committer for Parquet output. For Iceberg tables the table format handles commits, which is one more reason to convert.
Oozie workflows rarely have a managed equivalent; most teams rewrite them as Airflow DAGs or the target platform's workflow service. Port them in dependency order and dual-run: the job runs on both platforms against the same inputs for at least one full business cycle, and the outputs are compared automatically.
The security lane
Kerberos principals and Ranger policies do not map one to one onto cloud identity. The workable pattern is a role per workload rather than per person: each job runs as a service identity with access to the prefixes it needs, and analysts reach data through a catalog-level policy engine such as AWS Lake Formation, the target catalog's own grants, or Ranger on the managed platform. Export the Ranger policies, list who actually used each grant in the audit logs, and migrate only the grants that were exercised. HDFS encryption zones become bucket-level or prefix-level KMS keys, so decide key ownership before the first byte is copied.
Validation and cutover
Validation must compare content, not only sizes. The PySpark check below fingerprints each partition with a row count and an order-independent hash sum, then reports any partition that differs:
# Compare a table partition by partition: row counts plus an order-independent content hash
from pyspark.sql import functions as F
def fingerprint(df, keys):
row_hash = F.xxhash64(*[F.col(k) for k in df.columns])
return (df.groupBy(*keys)
.agg(F.count("*").alias("rows"),
F.sum(row_hash.cast("decimal(38,0)")).alias("h")))
src = fingerprint(spark.table("hive_onprem.sales.orders"), ["order_date"])
dst = fingerprint(spark.table("lake.sales.orders"), ["order_date"])
diff = (src.alias("s").join(dst.alias("d"), "order_date", "full_outer")
.where("s.rows IS NULL OR d.rows IS NULL OR s.rows != d.rows OR s.h != d.h"))
bad = diff.count()
print(f"mismatched partitions: {bad}")
assert bad == 0, diff.limit(20).collect()Cut over one table or one domain at a time. The sequence: freeze writers on-prem for the table, run a final incremental copy, run validation, switch the catalog entry and job configuration, and keep the on-prem copy read-only for a fixed rollback window, usually two to four weeks. Only after the window, with no rollback, delete the source.
Rollback must be rehearsed, not assumed. During the window, any table written in the cloud diverges from the frozen on-prem copy, so a rollback either replays the cloud writes back with DistCp or accepts re-running the jobs since cutover. Decide which, per table, before the switch.
Worked example: a wave plan
Take the 670 TB estate from the arithmetic above: 180 million files, 4,000 Hive tables, 900 Oozie workflows and a 10 Gbit/s link. The inventory shows 260 TB untouched for over a year, 40 tables holding 70 percent of the active bytes, and one log directory with 60 million files under 1 MB.
| Wave | Scope | Data | Exit criterion |
|---|---|---|---|
| 0, weeks 1-2 | Pilot: one medium partitioned table, its 6 jobs and its readers | 2 TB | All four lanes done, rollback rehearsed |
| 1, weeks 2-4 | Cold archives to infrequent-access storage; small-file log directory compacted, then copied | 260 TB | Spot-check restores; no on-prem readers in audit log |
| 2, weeks 4-8 | The 40 large active tables, converted to Iceberg; their jobs dual-run | 290 TB | Zero mismatched partitions for 5 business days |
| 3, weeks 8-14 | Long tail in waves of 200 tables, mostly registered as-is | 120 TB | Per-wave validation and job sign-off |
| 4, weeks 14-18 | Rollback windows close; DataNodes and NameNodes decommissioned | None | Source deleted, contracts ended |
The pilot matters more than its size suggests. It flushes out the missing pieces, such as a firewall rule, an IAM boundary or a Hive SerDe the new engine lacks, while the blast radius is one table.
Failure modes and trade-offs
- Copy succeeded, table broken. Partitions copied but never registered in the catalog are invisible to queries. Validate through the catalog, never through file listings.
- Writers during the copy. A job that writes to a partition after it was copied leaves the target stale. Freeze or dual-write before the final pass.
- Request cost and throttling. Millions of small objects produce request charges and per-prefix throttling. Compact first.
- Timezone and type drift. Older Hive timestamp semantics and newer engines can disagree on the same Parquet files; the content hash catches this when the comparison runs through both engines.
- Production impact. An uncapped DistCp saturates DataNode disks and the uplink. Run it in its own YARN queue with
-bandwidthlimits. - Shadow consumers. Tools reading HDFS directly, such as exports over WebHDFS, appear only in audit logs. Check the NameNode audit log for clients that are not in your job inventory.
What to do next
- Extract the fsimage and build a dataset inventory with logical size, file count, small-file count and idle days.
- Do the transfer arithmetic for your link and decide which cold data goes by appliance or not at all.
- Choose the target shape and the catalog, and decide which tables become Iceberg.
- Migrate one medium-sized partitioned table end to end through all four lanes, including validation and a rehearsed rollback.
- Turn that run into scripts: copy, register, validate and cut over, parameterised by table.
- Mine the NameNode audit log for unknown clients before freezing any writer.