Apache Hadoop is not dead. The project shipped 3.4.0 in March 2024 and 3.5.0 in April 2026, large on-premises clusters still run payroll, telecom billing and national statistics, and Hadoop's client libraries sit inside almost every Spark job that reads cloud storage. But the Hadoop platform, meaning HDFS for storage, YARN for scheduling and MapReduce or Hive on top, has stopped being the default place to build a data platform. New projects pick object storage, an open table format and engines that scale on their own.
Calling this "fashion" misses the engineering. Hadoop rested on a few design premises that were right for 2006 hardware and have since been undermined, one by one. This article goes through those premises and what replaced each, with the arithmetic that makes the argument concrete. It then covers what still legitimately belongs on Hadoop, a cost model for deciding whether your cluster should stay, and the ways migrations away from it go wrong. If you run a cluster today, the goal is to leave you able to make, and defend, the keep-or-leave decision.
The original premises
Hadoop copied the design of Google's GFS and MapReduce papers for clusters of cheap servers with local disks and, by today's standards, slow networks. Three premises followed. First, move computation to the data: since the network was the scarce resource, the scheduler places tasks on the node that holds the block, so most reads are local disk reads. Second, make failure normal: commodity disks die, so HDFS keeps three copies of every block and MapReduce re-runs failed tasks. Third, keep metadata simple: a single NameNode holds the whole namespace in memory, which keeps lookups fast and the design easy to reason about.
Each premise was a good bet at the time, and each produced a coupling. Locality ties storage and compute to the same machines. Replication triples raw capacity. In-memory metadata ties namespace size to one process's heap. The decline of Hadoop is mostly the story of those couplings becoming more expensive than what they bought.
Cause 1: storage and compute are welded together
On a classic cluster every worker is both a DataNode and a YARN NodeManager. When storage fills up you buy nodes, and you get CPU and memory you may not need. When a quarterly batch needs ten times the compute for six hours, you either own that capacity all year or the job waits. Typical clusters end up lopsided: disks full and CPUs mostly idle, or the reverse. Three-way replication means 1 PB of data needs about 3 PB of disk, plus free space for balancing and failures.
Object storage breaks the weld. You pay for logical bytes stored, and the provider handles durability. Compute engines start when needed and stop when not, on Kubernetes, on managed services such as EMR or Dataproc, or as serverless SQL. Once storage and compute are billed separately, idle capacity is a choice rather than a consequence of the design. This is the largest single economic cause of the decline.
Cause 2: the network caught up and object stores became safe
Data locality only pays when reading over the network is much slower than reading from local disk. Datacenter networks moved from 1 GbE to 25, 100 GbE and beyond, while columnar formats such as Parquet and ORC cut the bytes a query reads by skipping columns and row groups. A modern engine reading compressed columns from object storage over a fast network is rarely limited by the network, so the main reason to put compute next to disks has faded.
The second blocker was semantics. Early S3 was eventually consistent for some operations, so a job could list a directory and miss files just written, and Hadoop needed workarounds such as S3Guard. Since December 2020 S3 has offered strong read-after-write consistency, and the other major object stores also give strong consistency. One gap remains: object stores have no atomic directory rename, which the classic Hadoop output committer relied on. That gap was closed above the storage layer, by S3A committers and, more completely, by table formats that commit by atomically swapping a metadata pointer rather than renaming directories.
Cause 3: the NameNode does not scale with file counts
Every file, directory and block is an object in the NameNode's Java heap. The common rule of thumb is about 150 bytes each. Bytes do not matter; object counts do. The same petabyte stored as large files or as millions of small ones produces very different heaps:
# Rough NameNode heap for a namespace. ~150 bytes per object (file, directory or block)
# is the commonly quoted rule of thumb, not an exact figure; measure your own cluster.
BYTES_PER_OBJECT = 150
def namenode_heap_gib(files, dirs, avg_blocks_per_file):
objects = files + dirs + files * avg_blocks_per_file
return objects * BYTES_PER_OBJECT / 2**30
# Same 1 PiB of data, two layouts (128 MiB blocks, 3x replication does not add namespace objects)
print(namenode_heap_gib(files=8_000_000, dirs=200_000, avg_blocks_per_file=1.0)) # ~2.3 GiB
print(namenode_heap_gib(files=400_000_000, dirs=5_000_000, avg_blocks_per_file=1.0)) # ~113 GiBHundred-gigabyte heaps mean long garbage-collection pauses, slow start-up while the NameNode replays its edits and processes block reports, and slower failover. Federation and router-based federation split the namespace, at the price of more moving parts. The detailed mechanics are in the HDFS small files problem. Object stores spread metadata across a distributed service, so there is no single heap to outgrow (they have their own limits on request rates per prefix). Apache Ozone exists largely to give on-premises users that same property; see Apache Ozone in depth.
Cause 4: the compute layers were overtaken
MapReduce writes intermediate results to disk between every stage. Spark keeps them in memory where it can and plans whole multi-stage pipelines, which is why it replaced MapReduce for most batch work in the mid-2010s. SQL moved to engines such as Trino, Impala and cloud warehouses. MapReduce is still in the distribution, but few people write new MapReduce jobs.
YARN met the same fate from a different direction. It is a capable scheduler (see YARN in depth), but it schedules only JVM-centric data jobs, while organisations standardised on Kubernetes for everything else, including services, ML training and notebooks. Running two cluster managers means two sets of capacity, security and on-call. Spark gained a native Kubernetes scheduler backend, so the main tenant of YARN can move; see Spark on Kubernetes.
Finally, Hive's model of a table as a directory tree with partitions stored as subdirectories gave way to open table formats. Iceberg, Delta Lake and Hudi keep schema, partitions and snapshots in metadata files, which brings atomic commits, time travel, schema evolution and row-level deletes on object storage. Once the table format provides those, HDFS's semantics are no longer needed; see Apache Iceberg and Iceberg with Trino.
Cause 5: operations and the vendor shake-out
Running Hadoop well needs specialists: Kerberos and its keytabs, Ranger policies, rolling upgrades across HDFS, YARN, Hive, HBase and ZooKeeper with compatible versions, NameNode high availability with journal nodes, and balancers. Those skills are scarce and expensive, and they do not transfer to the rest of the platform.
The commercial market also consolidated. Cloudera and Hortonworks merged in January 2019, HPE acquired MapR's business later that year, and Cloudera was taken private in 2021. Customers who had planned around competing distributions faced licence changes and forced upgrade paths, which prompted many to re-evaluate rather than renew. The surviving vendor offerings now support object storage and Kubernetes themselves, which says a lot about where the architecture went.
What still belongs on Hadoop
- Large, steady, on-premises workloads. If the cluster runs hot around the clock and the hardware is paid for, elasticity buys little and egress-free local disks are cheap.
- Data that cannot leave a building. Sovereignty, defence and some regulated workloads keep data on owned hardware; Ozone or HDFS remain sensible stores there.
- HBase-dependent systems. HBase runs on HDFS, and replacing a low-latency store is a separate, larger project.
- Stable legacy pipelines near retirement. Rewriting a job that will be switched off in eighteen months rarely pays.
- The libraries. Even fully cloud-native stacks keep using hadoop-common, the S3A connector and the Hadoop configuration model, as below.
# The libraries outlive the platform: Spark on Kubernetes reading S3 still uses Hadoop's S3A.
spark-submit \
--conf spark.hadoop.fs.s3a.committer.name=magic \
--conf spark.sql.catalog.lake=org.apache.iceberg.spark.SparkCatalog \
--conf spark.sql.catalog.lake.type=rest \
--conf spark.sql.catalog.lake.uri=https://catalog.internal/api \
job.py s3a://raw-events/2026/10/01/
Worked example: should this cluster stay?
Take a 200-node cluster with about 12 TB usable per node: roughly 2.4 PB raw, about 800 TB of logical data after three-way replication. Monitoring shows the cluster busy about 35 percent of hours, with month-end peaks that still queue. Four engineers keep it running. The model below compares annual costs. Every price is a placeholder to replace with your own contracts and quotes; the structure is the point, not the totals.
# Keep-or-leave arithmetic. Every price below is a PLACEHOLDER: replace with your contracts.
def on_prem_annual(nodes, node_cost, years_amortised, power_space_per_node, ops_fte, fte_cost,
support_per_node):
return nodes * (node_cost / years_amortised + power_space_per_node + support_per_node) \
+ ops_fte * fte_cost
def cloud_annual(logical_tb, gb_month, compute_hours, hour_cost, requests_cost, ops_fte, fte_cost):
storage = logical_tb * 1024 * gb_month * 12 # object storage already replicates
return storage + compute_hours * hour_cost + requests_cost + ops_fte * fte_cost
onprem = on_prem_annual(nodes=200, node_cost=20_000, years_amortised=5,
power_space_per_node=2_500, ops_fte=4, fte_cost=180_000,
support_per_node=1_500)
# 200 nodes x 12 TB usable / 3 replicas = ~800 TB logical; cluster busy ~35% of hours
cloud = cloud_annual(logical_tb=800, gb_month=0.02, compute_hours=0.35 * 200 * 8760,
hour_cost=1.0, requests_cost=60_000, ops_fte=2, fte_cost=180_000)
print(f"on-prem ~${onprem:,.0f}/yr, cloud ~${cloud:,.0f}/yr")
# Then add one-off migration cost and the months of running both.Read the result honestly. Utilisation is the swing variable: at 35 percent busy, paying only for used compute hours usually wins; at 85 percent around the clock, owned hardware often wins. Request costs grow with file counts, so a small-files problem follows you to object storage as a bill. Staff cost moves rather than disappears, because platform work shifts to cost control, catalog governance and Kubernetes. Add the one-off migration (often a year of two systems running in parallel) and egress if data must come back. If the cloud number is not clearly lower after all of that, compaction, Ozone or a partial move (cold data to object storage, compute where it is) may be the better answer.
How migrations away from Hadoop fail
- Lift and shift of the namespace. Copying millions of small files to object storage multiplies request costs and slows listing. Compact before or during the copy.
- Rename-based job semantics. Jobs that write to a temporary directory and rename on success are not atomic on object stores. Move them to table-format commits or the S3A committers before switching storage.
- Hive tables carried over as directories. You keep the old problems on new storage. Convert to Iceberg or Delta during the move, and keep a catalog as the single source of table truth.
- Security rebuilt late. Kerberos principals and Ranger policies do not map one-to-one to IAM roles and catalog grants. Inventory who reads what first, then design the target model.
- Unowned long tail. The last 10 percent of jobs (cron scripts, ad hoc Pig, forgotten exports) keeps the old cluster alive for years. Find owners and decommission dates early.
- Double running without an end date. Parallel operation is necessary, but without per-dataset cut-over criteria it becomes permanent and the savings never arrive.
What to do next
- Measure before deciding: cluster utilisation by hour, storage growth, file and block counts, NameNode heap, and staff time spent on platform operations.
- Fill in the cost model with real contracts and quotes, including migration effort and a period of parallel running, and test it at your actual and peak utilisation.
- Fix what helps either way: compact small files, convert key Hive tables to an open table format, and move new Spark work to a scheduler that does not depend on YARN.
- If you leave, migrate by dataset, not by cluster: pick a domain, move its tables and jobs, prove parity with row counts and checksums, then cut readers over.
- If you stay, plan your upgrade path to the current 3.4.x or 3.5.x line, staff for it explicitly, and evaluate Ozone for namespace scale.
- Write the decision down with its assumptions and revisit it yearly; hardware refresh dates are the natural review points.