skip to content

When would you still build a new platform on HDFS instead of object storage?

level: principalimportance: nice to knowfreq 30%

answer

  1. separate the property from the era
  2. almost never, except on-premises
  3. who still has no object store?
  4. what replaces atomic rename today?
  5. the answer is a table format

basics

~20 s

Rarely, and only on-premises: an air-gapped or regulated cluster, or a latency-sensitive workload like HBase where compute sits on the same disks as the data. New cloud platforms default to object storage plus an open table format.

solid answer

~40 s

For a greenfield cloud platform the honest answer is almost never — object storage plus Iceberg, Delta Lake or Hudi has taken over, because it decouples capacity from compute, removes the NameNode as a namespace ceiling, and drops the operational burden of HA NameNodes, JournalNodes and ZooKeeper failover. HDFS still earns its place on-premises: an air-gapped or data-residency-constrained environment with no object store, a large existing Hadoop estate where migration cost outweighs the benefit, or a latency-sensitive layer such as HBase that wants local disks and short, predictable reads. Even then the modern on-prem answer is usually an S3-compatible layer like MinIO or Ceph rather than HDFS. What I would not do is choose HDFS for cheap bulk storage — three-way replication on always-on machines is the expensive option.

code

bash · 5 lines
bash
# Bulk-copy a Hive warehouse directory from HDFS to an S3A target
hadoop distcp \
  -update -delete \
  hdfs://nn1:8020/user/hive/warehouse/events \
  s3a://analytics-lake/warehouse/events

go deeper

for a junior

Know that HDFS stores files across a cluster of machines and that most new platforms instead keep data in cloud object storage such as S3, with engines reading the same Parquet files either way.

for a middle

Be able to name what you lose and gain in the swap: data locality and atomic rename on one side, independent scaling of storage and compute plus far less operational surface on the other.

for a senior

Demonstrate that you have run the operation — HA NameNode topology, namespace limits, rebalancing — and can sketch a per-dataset migration that also fixes the commit protocol and the file-size distribution.

for a principal

Own the decision and its blast radius: total cost at real utilisation, regulatory and residency constraints, which engines the org standardises on, and the fact that choosing a table format is a multi-year commitment for every team.

## The question behind the question This is a strategy question, not a trivia question. The interviewer wants to hear whether you can separate *what HDFS actually provides* from *the era it belongs to* — and whether you will give a defensible "almost never, except…" instead of either nostalgia or reflexive dismissal. ## What HDFS genuinely gives you **Data locality.** The NameNode knows which DataNodes hold each block, and schedulers use that: a Spark task can be placed on the machine already holding its input and read from local disk. On a well-provisioned on-prem cluster with modest network bandwidth, this is a real throughput advantage. **A real filesystem contract.** A hierarchical namespace, atomic directory rename, append and truncate. Anything that was written against filesystem semantics — the classic `FileOutputCommitter`, Hive's dynamic-partition insert, older ecosystem tools — works without modification. **Predictable low-latency reads on local disks.** HBase is the standing example: a random-read store that wants short, consistent I/O paths rather than per-request object-store latency. **It runs where there is no cloud.** Air-gapped defence and government environments, strict data-residency regimes, or a site with no object store at all. ## What it costs **Operational weight.** A production HDFS is not one service. Highly available deployments need an active and a standby NameNode, a JournalNode quorum for the shared edit log, and ZooKeeper plus failover controllers to arbitrate. Someone has to own upgrades, rebalancing, decommissioning and the occasional metadata recovery. **A namespace ceiling.** Every file, directory and block is an object in the NameNode's heap, so the metadata capacity of a cluster is bounded by one JVM's memory. Federation and router-based federation exist to shard the namespace, but they add topology rather than remove the constraint. **Coupled scaling.** Capacity comes from DataNodes, and DataNodes are machines you keep powered whether or not you are computing. You cannot grow a petabyte of cold storage without also growing the compute footprint. **Storage multiplier.** Default replication is three copies. Hadoop 3's erasure coding cuts that materially for cold data, at the cost of more expensive reconstruction on failure, but it is still local disk you own. ## The modern default, and what replaces each property A greenfield platform in 2026 is object storage (S3, ADLS, GCS; MinIO or Ceph on-prem) holding columnar files, with an open table format — Iceberg, Delta Lake or Hudi — layered on top, and a catalog for discovery. The table format is what replaces the property people actually relied on HDFS for: **atomic commit**. Instead of renaming a directory into place, the engine writes new data files and then atomically swaps a metadata pointer, which also brings snapshot isolation, time travel and schema evolution that HDFS never offered. Locality is replaced by aggregate bandwidth plus aggressive pruning: large files, partitioning, predicate and projection pushdown, and caching. Multiple engines can read the same tables concurrently without sharing a cluster. ## So when does HDFS still win Be concrete. On-premises and air-gapped, with no S3-compatible layer available. An existing Hadoop estate where the data is already there and migration cost dominates. A co-located HBase or similar latency-sensitive store. A workload whose economics genuinely favour owned hardware at sustained high utilisation. Notice that every one of these is an existing-context argument, not a design-from-scratch one — which is exactly the point. ## The migration angle If the answer is "move off", say how. Bulk-copy with `hadoop distcp` to an `s3a://` target, re-register tables in the new catalog, convert to a table format, and — the step people forget — fix the commit path, because a job that relied on rename-based commit will be slow or unsafe against object storage until it uses an S3A committer or a table format. Run both in parallel and cut over per dataset rather than in one move. One clarification worth volunteering: Spark's shuffle never used HDFS. Intermediate shuffle data goes to executor-local disks under `spark.local.dir`, so "we need HDFS for shuffle" is not a reason to keep it.

  • What breaks first when you point an existing Hadoop workload at object storage unchanged?
    The commit path. Anything relying on rename — the classic FileOutputCommitter, Hive dynamic-partition inserts — turns a metadata edit into a full copy of the output, so the final stage finishes fast and the job then sits in commit. Directory listings over many small objects are the second casualty. Fix both with an S3A committer or a table format, plus larger files.
  • Would you keep an HDFS cluster just to hold Spark's shuffle data?
    No — Spark never wrote shuffle data to HDFS. Shuffle blocks go to executor-local disks configured by `spark.local.dir`, optionally served by an external shuffle service. If shuffle I/O is the bottleneck the answer is faster local disks, better partitioning or less shuffling, not a distributed filesystem.
  • If you must stay on-premises, why choose MinIO or Ceph over HDFS?
    Because you get the S3 API, which is what every current engine and table format targets, without a NameNode namespace ceiling or a JournalNode-plus-ZooKeeper HA topology. You also decouple capacity from compute nodes. You give up data locality, which matters far less once reads are pruned and files are large.

saying these in an interview costs you the question

  • Says HDFS is the cheap option for cold bulk storage
  • Recommends HDFS for a greenfield cloud platform with no on-prem constraint
  • Claims object storage cannot support ACID table updates
  • Believes Spark requires HDFS to run
  • Treats migration as a copy, ignoring the commit path and file sizes

context