"Should we still be on Hadoop?" is usually argued from opinion: one side points at cloud object stores, the other at a paid-off rack. Both sides are often right about different workloads on the same cluster. This guide treats the question as a measurement problem. You will collect three kinds of evidence from the cluster itself, score each workload separately, and end with a placement for each one: stay, hybrid or move.
The history of why Hadoop lost its default status, and a keep-or-leave cost model, are covered in Why Hadoop declined; this page does not repeat them. It is about producing the inputs those models need, and about what staying responsibly looks like if the answer for some workloads is stay.
Why the answer is per workload
A typical long-lived cluster hosts a mix: a nightly ETL that saturates the cluster for six hours, an HBase table serving lookups, a few hundred ad hoc Hive queries, and data that nobody has read in two years. These have different economics. The ETL uses owned hardware well. The ad hoc queries are spiky and would cost less on elastic compute. The cold data costs disks, power and NameNode memory for no reads at all.
Deciding for the whole cluster forces one answer onto all of them. Deciding per workload lets you move what benefits from moving, keep what benefits from staying, and shrink the cluster to the part that earns its place. Figure 1 shows the flow.
What Hadoop is still mechanically good at
Hadoop's original advantage was moving compute to the data. YARN schedules containers on the nodes that hold the HDFS blocks they read, and short-circuit local reads let a task read a block file straight from local disk, bypassing the DataNode's network path. For a large job that scans most of its input, local disks deliver aggregate bandwidth that grows with node count, and you pay no per-request or per-gigabyte fees.
The second advantage is cost shape. Owned hardware is a fixed cost: once bought, an hour of use costs power and staff time, not a metered price. A cluster that is busy most of the day spreads that fixed cost thinly. A cluster that is idle most of the day is paying for capacity it does not use, which is exactly what elastic compute avoids.
Neither advantage is universal. Locality matters less when the network is fast and the job is compute-bound. Fixed cost is only cheap if utilisation is high. So the evidence you need is utilisation over time, the shape of your data, and the constraints you cannot change.
Evidence 1: utilisation over time
The ResourceManager exposes cluster totals at /ws/v1/cluster/metrics: allocated and total memory and vcores, plus pending applications. Sampling it every five minutes for a month gives you the utilisation curve that every cost argument depends on. A single snapshot is worthless; the shape is the point.
# Sample YARN allocation every 5 minutes for 30 days; one JSON line per sample.
import json, time, urllib.request
RM = "http://rm.internal:8088/ws/v1/cluster/metrics" # add SPNEGO auth on Kerberised clusters
def sample():
m = json.load(urllib.request.urlopen(RM, timeout=10))["clusterMetrics"]
return {
"ts": int(time.time()),
"mem_util": m["allocatedMB"] / max(m["totalMB"], 1),
"vcore_util": m["allocatedVirtualCores"] / max(m["totalVirtualCores"], 1),
"apps_pending": m["appsPending"],
}
with open("rm_samples.jsonl", "a") as out:
while True:
out.write(json.dumps(sample()) + "\n"); out.flush()
time.sleep(300)import json, statistics
s = [json.loads(l) for l in open("rm_samples.jsonl")]
u = sorted(max(x["mem_util"], x["vcore_util"]) for x in s)
pct = lambda q: u[int(q * (len(u) - 1))]
busy = sum(1 for x in u if x > 0.6) / len(u)
queued = sum(1 for x in s if x["apps_pending"] > 0) / len(s)
print(f"p10 {pct(.10):.0%} median {pct(.5):.0%} p90 {pct(.9):.0%}")
print(f"hours above 60% allocated: {busy:.0%}; hours with queued apps: {queued:.0%}")Read three things from the output. The median says whether owned hardware is working most of the time. The ratio of p90 to median says how spiky demand is; a ratio above about three means you are sizing hardware for peaks that rarely happen. The share of samples with queued applications says whether the cluster is already too small at peak, which is the case where bursting some work elsewhere helps even if you stay.
To get per-workload numbers, sample /ws/v1/cluster/scheduler for per-queue usage too, or aggregate finished applications from /ws/v1/cluster/apps by queue or user. Queues usually map well to teams and workloads.
Evidence 2: what the data looks like
The NameNode keeps every file, directory and block in memory; a common rule of thumb is on the order of 150 bytes per object, so object count rather than terabytes is what limits it. File-size distribution also predicts the cost of moving, because object stores bill per request and many small files mean many requests. Profile the namespace offline from a checkpoint, so the audit puts no load on the NameNode:
# On a gateway node: fetch the latest checkpoint, then dump it offline (no NameNode load).
hdfs dfsadmin -fetchImage /data/audit/
hdfs oiv -p Delimited -i /data/audit/fsimage_0000000000123456789 \
-o /data/audit/fsimage.tsv# Columns include Path, Replication, ..., BlocksCount, FileSize, ..., Permission.
import csv, collections
buckets = collections.Counter(); bytes_by = collections.Counter()
edges = [(1 << 20, "<1MB"), (16 << 20, "1-16MB"), (128 << 20, "16-128MB"), (1 << 40, ">128MB")]
with open("fsimage.tsv") as f:
for row in csv.DictReader(f, delimiter="\t"):
if row["Permission"].startswith("d"):
continue # directory, not a file
size = int(row["FileSize"])
top = "/".join(row["Path"].split("/")[:3]) # e.g. /warehouse/clicks
label = next(l for edge, l in edges if size < edge)
buckets[(top, label)] += 1; bytes_by[top] += size
for (top, label), n in sorted(buckets.items()):
print(f"{top:30} {label:9} {n:>12,}")Group by the top two path levels, which usually correspond to databases or pipelines. Add access time if dfs.namenode.accesstime.precision is enabled on your cluster: paths with no reads in a year are candidates for erasure coding, archival or deletion regardless of the platform decision. A directory with tens of millions of sub-megabyte files is a problem on either platform; The HDFS small files problem covers compaction, which you should do before a move rather than paying per-request fees for the mess afterwards.
Evidence 3: constraints and gravity
Some facts end the discussion for a workload before any arithmetic:
- Residency. Data that contracts, regulation or an air gap keep on owned hardware stays. The question becomes which on-premises store, HDFS or Apache Ozone, not whether to move.
- HBase. HBase stores its files in HDFS. Moving it means replacing a low-latency serving store, a separate project with its own risks.
- Gravity. If the consumers of a dataset are on premises, moving the data creates egress traffic on every read. Count who reads it, from where, and how much.
- Retirement. A pipeline due to be switched off within about eighteen months rarely repays a rewrite.
- People. If fewer than two people can operate HDFS and YARN, the on-call rota is a risk in itself, whatever the hardware costs.
A per-workload scorecard
Turn the evidence into a repeatable decision. The function below is deliberately simple; the weights are a starting point to argue about in a review, not a model of your costs. What matters is that every workload is scored on the same measured facts and the reasons are written down.
def placement(w):
"""w: dict of measured facts for ONE workload. Returns stay / hybrid / move and the reasons."""
if w["must_stay_on_prem"]: # residency, contract, air gap
return "stay", ["hard constraint"]
stay, move, why = 0, 0, []
if w["median_util"] >= 0.6: stay += 2; why.append("steady high utilisation")
if w["p90_util"] / max(w["median_util"], 0.01) > 3: move += 2; why.append("spiky demand")
if w["small_file_share"] > 0.5: move += 1; why.append("small files (fix before or after a move)")
if w["reads_per_tb_day"] > 50: stay += 1; why.append("hot local reads")
if w["depends_on_hbase"]: stay += 2; why.append("HBase dependency")
if w["retire_within_months"] <= 18: stay += 2; why.append("retiring soon: do not rewrite")
if w["ops_staff_left"] < 2: move += 2; why.append("cannot staff on-call")
if abs(stay - move) <= 1:
return "hybrid", why
return ("stay" if stay > move else "move"), whyFeed the result into the cost comparison in Why Hadoop declined for each workload that scores move or hybrid. A score is not a budget; it tells you which workloads are worth costing in detail.
Worked example: one cluster, three answers
Consider a 120-node cluster. A month of sampling shows a median allocation of 48 percent and a p90 of 95 percent, with queued applications in 12 percent of samples, all between 01:00 and 07:00. Its three main workloads score as follows.
| Workload | Evidence | Placement |
|---|---|---|
| Nightly billing ETL | Runs 01:00-07:00 at 90%+; reads 400 TB of local Parquet; billing data must stay in-country | Stay |
| Analyst Hive queries | Daytime only; p90/median about 5; 30% of files under 1 MB | Move, after compaction |
| Clickstream archive | 1.1 PB, 85% of paths unread for a year | Hybrid: erasure-code now, export cold partitions later |
The outcome is not "leave Hadoop" or "stay on Hadoop". It is a smaller cluster sized for the billing ETL, analyst queries moved to an elastic engine reading compacted tables, and the archive shrunk in place. The overnight queueing also disappears, because the analysts' long-tail jobs were the ones overlapping the ETL window.
Staying well: the baseline
Workloads that stay need a maintained platform. A cluster frozen on an old release is how "Hadoop still makes sense" turns into "Hadoop is stuck". The minimum baseline:
- A supported release line. Hadoop 3.5.0 requires Java 17 on the server side and supports Java 17 and 21 for clients, so plan the JDK upgrade with the Hadoop upgrade. The same release removed the deprecated WASB connector in favour of ABFS and added a Google Cloud Storage FileSystem. The Hadoop ecosystem in 2026 maps which surrounding projects are active and their Java baselines.
- Open table formats. Writing Iceberg tables on HDFS, with a catalog that other engines can reach, keeps the exit open: the same tables can later be copied to object storage and read by other engines without rewriting pipelines.
- Erasure coding for cold data. HDFS erasure coding with the
RS-6-3-1024kpolicy stores 1.5 bytes per logical byte instead of 3 for three-way replication, at the cost of slower reconstruction and reads. Apply it to cold paths only. - Burst through connectors. The S3A and ABFS connectors let jobs on the cluster read or write object storage, so peak work or new projects can land elsewhere without a migration.
- Security as standard. Kerberos, authorisation policies and encryption are table stakes when the cluster holds the data that could not move.
Failure modes
- Deciding from a snapshot. One busy afternoon makes any cluster look full. Sample for at least a month, including a month-end.
- Whole-cluster verdicts. Moving everything drags hard-constrained workloads into exceptions; keeping everything keeps paying for idle capacity.
- Moving the small files. Copying millions of tiny files to object storage turns a NameNode problem into a request bill. Compact first.
- Frozen in place. Staying without upgrades accumulates JDK, security and connector debt until a forced migration happens on someone else's schedule.
- Hybrid without ownership. A split platform where nobody owns the boundary ends with data copied both ways and nobody sure which copy is authoritative. Assign each dataset one home.
Trade-offs
Staying keeps locality, fixed costs and data under your control, and costs you hardware refreshes, upgrades and scarce specialist staff. Moving buys elasticity and managed operations, and costs migration effort, request and egress fees, and a period of running both platforms; The Hadoop cloud migration playbook covers that path. Hybrid gets the best placement per workload at the price of two platforms to run. Whichever you pick, open table formats and measured evidence keep the decision reversible.
What to do next
- Start sampling ResourceManager metrics every five minutes today; you need a month of data.
- Fetch a checkpoint, dump it with hdfs oiv, and produce file-count and size distributions per top-level path.
- List hard constraints per workload: residency, HBase, consumers' locations, retirement dates, staffing.
- Score each workload with the same function and record the reasons in a decision document.
- Cost only the workloads that score move or hybrid, using your contracts rather than list prices.
- For what stays, schedule the move to a supported release and Java 17, start writing new tables as Iceberg, and erasure-code cold paths.