HDFS, the Hadoop Distributed File System, stores very large files across many commodity servers and keeps them available when disks and machines fail. Its design comes from Google's 2003 GFS paper and makes a set of deliberate bets: files are large, written once and read many times, read in long sequential scans, and processed by computation that can move to the data. Those bets make HDFS excellent at feeding Spark, Hive and MapReduce and poor at the workloads that break them, such as millions of small files or random updates.
This article explains the architecture from first principles, follows a write and a read through the system, shows how replication survives failures, works through the arithmetic of storing a real dataset, and gives the commands and trade-offs you need to decide whether HDFS or object storage fits your workload. Each component has a deeper article linked along the way.
The design assumptions
Every HDFS design choice follows from a few assumptions. Hardware fails constantly. In a cluster of thousands of disks some are always dead, so failure handling is routine work, not an exception path. Throughput beats latency. A batch job scanning a terabyte cares about aggregate megabytes per second, not the 10 ms to open a file. Write once, read many. Files are created, written and closed; afterwards you can append but not modify bytes in the middle, which removes most of the hard consistency problems. Moving computation is cheaper than moving data. The filesystem tells schedulers where blocks live so tasks can run on the same node or rack.
Relaxing POSIX is the price. There is no random write, a single writer per file is enforced by leases, and metadata operations go through one service. In exchange, HDFS scales to petabytes on ordinary servers with local disks.
The design assumptions
Every HDFS design choice follows from a few assumptions. Hardware fails constantly. In a cluster of thousands of disks some are always dead, so failure handling is routine work, not an exception path. Throughput beats latency. A batch job scanning a terabyte cares about aggregate megabytes per second, not the 10 ms to open a file. Write once, read many. Files are created, written and closed; afterwards you can append but not modify bytes in the middle, which removes most of the hard consistency problems. Moving computation is cheaper than moving data. The filesystem tells schedulers where blocks live so tasks can run on the same node or rack.
Relaxing POSIX is the price. There is no random write, a single writer per file is enforced by leases, and metadata operations go through one service. In exchange, HDFS scales to petabytes on ordinary servers with local disks.
The architecture
An HDFS cluster has one active NameNode and many DataNodes. The NameNode holds the namespace (directories, files, permissions) and the mapping from each file to its blocks and from each block to the DataNodes holding replicas. It stores no file data. Its internals are covered in the NameNode deep dive.
DataNodes run one per storage server. Each stores blocks as ordinary files on its local disks, serves reads and writes, heartbeats to the NameNode every three seconds and reports which blocks it holds. See the DataNode deep dive for disk layout and reporting.
Production clusters add a standby NameNode that replays the active's edit log from a quorum of JournalNodes, with ZooKeeper-based failover controllers to promote it, as explained in HDFS high availability. The client is a library inside every application. It asks the NameNode where blocks are and then talks to DataNodes directly, so file bytes never pass through the metadata server.
Files, blocks and replicas
A file is split into fixed-size blocks, 128 MB by default (dfs.blocksize), and each block is replicated, three times by default (dfs.replication). The last block of a file is only as large as the remaining data; a 1 KB file uses 1 KB of disk per replica, not 128 MB. The cost of a small file is metadata, not space: every file and block is an object in NameNode memory.
Large blocks keep the NameNode's object count low and make a block read mostly sequential disk transfer rather than seeks. They also set the unit of parallelism: a Spark job typically gets one input split per block, so a 10 GB file yields about 80 tasks. Replication can be set per file, for example higher for a small lookup table read by every task.
Placement is rack-aware. With default policy the first replica goes on the writer's node if it is a DataNode (otherwise a random node), the second on a node in a different rack, and the third on another node in that second rack. One rack can fail without losing data, while only one replica crosses the inter-rack network during writes. Blocks and replication covers the policy in detail, and erasure coding covers the alternative to triple replication for cold data.
The write path
Writing a file, step by step. The client calls create; the NameNode checks permissions, adds the file entry and grants the client a lease, the lock that makes it the only writer. For each block the client asks for a new block, and the NameNode returns a list of DataNodes, one per replica, ordered as a pipeline. The client cuts the data into packets of about 64 KB with a checksum for every 512 bytes and sends them to the first DataNode, which stores and forwards to the second, which forwards to the third. Acknowledgements travel back up the pipeline. When the last block is done the client calls complete.
If a DataNode in the pipeline fails mid-block, the client removes it, the surviving replicas agree on a new generation stamp, and writing continues with the remaining nodes; depending on the replace-datanode-on-failure policy, a replacement may be added. The NameNode re-replicates the block later. In application code the whole protocol hides behind a stream:
Configuration conf = new Configuration(); // reads core-site.xml, hdfs-site.xml
FileSystem fs = FileSystem.get(URI.create("hdfs://nameservice1"), conf);
Path out = new Path("/data/events/2026-10-04/part-0000.parquet");
try (FSDataOutputStream os = fs.create(out, true)) { // create RPC, lease granted
os.write(bytes); // packets stream down the pipeline
os.hflush(); // visible to new readers, not fsynced
} // close() -> complete RPC
for (BlockLocation b : fs.getFileBlockLocations(fs.getFileStatus(out), 0, Long.MAX_VALUE)) {
System.out.println(b.getOffset() + " " + String.join(",", b.getHosts()));
}hflush makes written data visible to readers; hsync additionally asks DataNodes to sync to disk. Both matter for logs and write-ahead files, less for batch output.
The read path
To read, the client calls open and receives the locations of the first blocks, with replicas sorted by network distance from the reader: same node, then same rack, then elsewhere. It reads each block from the closest replica and verifies checksums as bytes arrive. On a checksum mismatch or a dead node it switches to another replica and reports the corrupt one to the NameNode, which schedules a fresh copy from a good replica. Readers never see corrupted bytes unless every replica is bad.
When the reader runs on the DataNode that holds the block, short-circuit local reads let the client open the block file directly through a UNIX domain socket handoff, skipping the DataNode's TCP path. Enable it with dfs.client.read.shortcircuit=true and a dfs.domain.socket.path; it is one of the larger wins for co-located Spark or HBase. Because YARN schedules tasks where blocks live (see the YARN overview), most reads in a busy cluster are local.
How HDFS survives failures
Failure handling is mostly automatic. A failed disk takes its blocks offline; the DataNode keeps running with its other disks if dfs.datanode.failed.volumes.tolerated allows, and reports the loss. A dead DataNode is declared after about ten and a half minutes without heartbeats, and the NameNode copies its blocks from surviving replicas to other nodes, most-endangered blocks first. Silent corruption is caught by checksums on read and by a periodic block scanner on each DataNode.
The NameNode is the component that needs explicit design. Without HA, its failure stops the cluster even though no data is lost. With HA, failover to the standby takes seconds to tens of seconds; clients retry through a logical nameservice URI such as hdfs://nameservice1. A Secondary NameNode, despite the name, is not a failover target: it only performs checkpoints in non-HA clusters.
Planned removal is gentler than failure. Add a host to the exclude file named by dfs.hosts.exclude and run hdfs dfsadmin -refreshNodes; the node enters decommissioning, its blocks are copied elsewhere while it keeps serving reads, and it is marked decommissioned only when every block has enough replicas on other nodes. Shutting a node off without this step works, but leaves the cluster under-replicated for the detection delay plus the copy time.
Worked example: a year of event data
Plan storage for 1 TB of new event data per day kept for one year. Raw replicated footprint: 365 TB of logical data times three replicas is about 1.1 PB, before free-space headroom; clusters usually run below about 75 to 80 percent full so re-replication and balancing have room, which pushes raw capacity toward 1.4 PB. Moving data older than 90 days to Reed-Solomon 6+3 erasure coding stores 1.5 bytes per logical byte instead of 3: 275 TB of cold data drops from 825 TB to about 412 TB raw, saving roughly 400 TB.
Now the NameNode. Written as 128 MB blocks in files of about 1 GB, a day is around 1,000 files and 8,000 blocks, so a year is about 3 million blocks: trivial. Written as 64 KB files by a careless stream job, the same terabyte is about 16 million files and 16 million blocks per day, over 10 billion objects a year, far beyond any single NameNode heap. The bytes are identical; file size decides whether the cluster survives. Compact streaming output into large files before it lands, as discussed in the NameNode article.
Network matters too. Each logical terabyte written with three replicas lands as three terabytes on disk and crosses the network at least twice, since the client usually writes the first replica locally and the pipeline forwards two copies, one of them between racks. Averaged over a day that is modest, around 12 MB/s of logical ingest, but ingest is rarely spread evenly: a two-hour nightly load needs twelve times that, and the inter-rack links carry a full copy of it.
Everyday operations
The handful of commands that cover daily operation:
hdfs dfs -ls /data/events # list; -du -h for sizes
hdfs dfs -put local.parquet /data/events/ # upload
hdfs dfs -setrep -w 5 /ref/lookup.csv # change replication, wait for it
hdfs fsck /data/events -files -blocks -locations # block health and placement
hdfs dfsadmin -report # capacity, live and dead DataNodes
hdfs dfsadmin -safemode get
hdfs ec -setPolicy -path /archive -policy RS-6-3-1024kMonitor capacity used per DataNode and cluster-wide, missing and under-replicated blocks (missing should always be zero), dead DataNodes, NameNode heap and RPC latency. fsck reads metadata only and is safe on a live cluster, though slow on very large trees.
Failure modes
| Failure | Cause | Prevention |
|---|---|---|
| NameNode out of memory | Too many small files | Compaction, file-count quotas, monitoring FilesTotal |
| Cluster full | Retention not enforced, trash not expiring | Capacity alerts at 70 percent, lifecycle jobs |
| Missing blocks | Several nodes lost before re-replication finished | Rack-aware placement, replace failed disks promptly |
| Write failures | Too few DataNodes with space for the pipeline | Headroom; balance after adding nodes |
| Cluster-wide stall | NameNode GC pause or lock contention | Heap sizing, FairCallQueue, HA |
Trade-offs
HDFS versus object storage. S3, GCS and Azure Blob separate storage from compute, scale without a NameNode, and charge only for bytes stored. HDFS offers data locality, atomic directory rename (which some commit protocols rely on), consistent low-latency metadata and predictable cost on owned hardware. In the cloud, object storage usually wins; on premises with steady heavy scans, HDFS often does. Replication versus erasure coding. Three replicas cost 200 percent overhead but give fast local reads and cheap recovery; RS 6+3 costs 50 percent and a normal read touches the six data units on six nodes, parity only during reconstruction, which makes recovery CPU- and network-heavy, so use it for cold data. Block size. Bigger blocks mean fewer objects and fewer tasks; smaller ones mean more parallelism on small datasets. The default suits most work.
What to do next
- Run
hdfs dfsadmin -reportandhdfs fsck /on your cluster and record missing, under-replicated and dead-node counts. - Compute average file size per top-level directory; anything far below the block size is a compaction target.
- Confirm NameNode HA is configured and test a manual failover in a maintenance window.
- Enable short-circuit local reads if compute runs on DataNodes.
- Identify data older than 90 days and estimate the savings of an erasure-coding policy.
- Read the NameNode and DataNode deep dives next, then replication and HA.