skip to content

HDFS

Distributed file storage for big data: files split into large replicated blocks, a NameNode holding the namespace, DataNodes holding the bytes. HDFS is named on its own in plenty of job descriptions and architecture diagrams, and the detailed ground is covered under the Hadoop tree.

questions

2

How does HDFS differ from cloud object storage like S3 for a Spark job's output?

level: middleimportance: should knowfreq 42%

answer

  1. one is a tree, one is a flat keyspace
  2. the expensive part is the commit
  3. rename is where the two diverge
  4. copy-then-delete is neither cheap nor atomic
  5. multipart upload replaces the rename

basics

~20 s

HDFS is a hierarchical filesystem where directory rename is an atomic metadata operation, so committing output is nearly free. S3 is a flat key store with no rename: a commit copies every object, so Spark needs a dedicated committer.

solid answer

~50 s

HDFS gives you a real filesystem: directories exist, and `rename` is an atomic edit to the NameNode's namespace, independent of data size. Hadoop's classic `FileOutputCommitter` relies on exactly that — tasks write to a temporary path and commit by renaming into place. Object storage has no directories and no rename; a "rename" on S3 is a server-side copy of every object plus a delete, so commit time scales with output size and is not atomic. Modern setups avoid the problem with the S3A committers (`fs.s3a.committer.name` = `directory`, `partitioned` or `magic`), which build the output through multipart uploads and make the commit a completion call rather than a copy. The other differences are locality — Spark can schedule a task on the node holding an HDFS block, whereas object storage is always over the network — and scaling: object storage decouples capacity from compute.

code

properties · 6 lines
properties
# HDFS destination: commit is a rename, effectively free
spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version=1

# S3A destination: commit via multipart upload instead of rename
spark.hadoop.fs.s3a.committer.name=magic
spark.sql.sources.commitProtocolClass=org.apache.spark.internal.io.cloud.PathOutputCommitProtocol

go deeper

for a junior

Know that HDFS is a filesystem with directories and object storage is a flat bucket of keys, and that both can hold the same Parquet files a Spark job reads and writes.

for a middle

Be ready to explain why job commit is cheap on a filesystem with atomic rename and expensive on a key store that has to copy objects, and to name a committer or table format that solves it.

for a senior

Show you have operated this: diagnosing a job whose final stage finishes in seconds but whose commit runs for twenty minutes, choosing a committer, and reshaping output file sizes for listing and read cost.

for a principal

Own the platform-level consequence — decoupling storage from compute changes cluster sizing, cost model and failure domains, and pushes correctness guarantees from the filesystem up into a table format you must then standardise on.

## Two different storage contracts HDFS and S3-style object storage both hold large files for a Spark job, but they expose different contracts, and almost every practical difference falls out of that. **HDFS** is a POSIX-*like* distributed filesystem. It has a genuine hierarchical namespace held in the NameNode's memory: directories are first-class objects, paths are entries in a tree, and metadata operations mutate that tree. Files are write-once with append supported; you cannot overwrite an arbitrary byte range in the middle of a file. The bytes themselves live as replicated blocks on DataNodes. **Object storage** (S3, ADLS, GCS, or on-prem MinIO/Ceph) is a flat key–value store. `s3a://bucket/warehouse/events/dt=2026-01-01/part-0.parquet` is one opaque key; the slashes are convention, not structure. There is no directory object, no atomic multi-key operation, and no rename primitive. ## Why the commit protocol is the real difference Spark does not let tasks write straight to the final path, because a task can fail or be speculatively duplicated. Hadoop's `FileOutputCommitter` therefore has tasks write to a temporary location and *rename* into the destination when they commit; the algorithm is selectable with `mapreduce.fileoutputcommitter.algorithm.version` (version 1 renames task output into a job-temp directory and does a final job-level rename; version 2 promotes each task's output directly, which is faster but exposes partial output if the job dies). On HDFS, that rename is a single edit in the NameNode's namespace: constant time, atomic, no data movement. On S3, the same call is emulated by the connector as *copy every object, then delete the originals*. A 500 GB output means 500 GB of server-side copying at commit time, in a step that is neither atomic nor cheap — and with algorithm version 1 the data is effectively moved twice. The fix is not to configure S3 differently but to use a committer designed for it. The S3A committers (`fs.s3a.committer.name` set to `directory`, `partitioned` or `magic`) exploit S3's multipart upload: task output is uploaded but the multipart upload is left *uncompleted*, and job commit simply completes the uploads for the tasks that succeeded and aborts the rest. Spark wires this up through `spark.sql.sources.commitProtocolClass`. Table formats such as Iceberg, Delta Lake and Hudi sidestep the issue entirely by committing a metadata pointer rather than moving files. ## Consistency Older advice says S3 is eventually consistent and therefore unsafe for job output — this is out of date. S3 has provided strong read-after-write consistency for all operations since December 2020, and the old S3Guard-style workarounds are gone. What S3 still does *not* give you is atomic multi-object commit, which is a different property and the one that actually matters here. ## Locality and throughput HDFS was designed so compute runs where the bytes are. The NameNode reports block locations, and Spark's scheduler uses them: it will wait (`spark.locality.wait`) to place a task at `NODE_LOCAL` rather than fall back to `RACK_LOCAL` or `ANY`. With object storage every read is remote and locality levels collapse to `ANY`. In practice this matters less than it sounds. Cloud object stores deliver very high aggregate bandwidth across many parallel range GETs, and columnar formats with predicate and projection pushdown mean you fetch a fraction of the bytes. What you must respect instead is per-request latency and listing cost: many small objects are expensive to list and read, so you target large files and partition prune aggressively. ## Scaling, durability and cost HDFS couples storage to compute — more capacity means more DataNodes, which you also pay to keep running — and its namespace is bounded by NameNode heap. Durability comes from replication (`dfs.replication`), with Hadoop 3 erasure coding available to cut the storage multiplier. Object storage decouples the two: you scale capacity and compute independently, pay per gigabyte stored, and get durability from the provider. ## What to say in an interview Lead with the contract difference — filesystem with atomic rename versus flat key store without one — then show the consequence (commit protocol), then the secondary axes (locality, scaling, listing cost). Naming the S3A committers, or Iceberg/Delta metadata commits, is what distinguishes someone who has actually shipped a job to object storage.

  • Why do the S3A committers use multipart upload instead of writing files and renaming them?
    Because a multipart upload can be started, fully uploaded, and left uncompleted. Task output is streamed to its final key but is invisible until the upload is completed, so job commit becomes a small set of complete calls for successful tasks and aborts for the rest. No bytes are copied at commit time, and partial output never becomes visible.
  • Does Spark still need data locality when it reads from object storage?
    No — every task reads over the network, so scheduler locality collapses to ANY and `spark.locality.wait` stops helping. You compensate on the data-layout side instead: large columnar files, partition pruning, predicate and projection pushdown so you fetch far fewer bytes, and caching hot datasets in executor memory.
  • Which HDFS operations does object storage not support at all?
    There is no rename and no atomic multi-key operation, so no atomic directory swap. There are no real directories, so listing is a prefix scan whose cost grows with object count. There is also no append or truncate in the HDFS sense, and no block-location reporting for locality-aware scheduling.

saying these in an interview costs you the question

  • Says S3 is eventually consistent so job output may be lost
  • Claims S3 splits files into blocks and replicates them like HDFS
  • Describes an S3 rename as a cheap metadata operation
  • Thinks Spark writes its shuffle data to HDFS
  • Assumes Spark cannot perform on object storage without data locality

context

open as a page

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

level: principalimportance: nice to knowfreq 30%

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.

open as a page