A Hadoop cluster is two distributed systems sharing the same machines. HDFS stores bytes: it splits files into large blocks, replicates them across racks and keeps the namespace in one process's memory. YARN allocates compute: it carves each node's memory and cores into containers and hands them to applications such as Spark, MapReduce, Tez or Flink. Neither knows much about the other's internals, yet the whole value of the platform comes from their cooperation, because YARN tries to run each task on a node that already holds the block it reads.
This article explains both systems from first principles, follows one job from submission to commit, sizes a real cluster with arithmetic you can reuse, and lists the failure modes that actually take Hadoop clusters down. It assumes you know what a file system and a scheduler are, and nothing else about Hadoop. Where the article quotes a default, it names the configuration key so you can check your own distribution, because vendors and versions do change defaults.
HDFS: blocks, NameNode and DataNodes
Blocks. HDFS stores every file as a sequence of fixed-size blocks, 128 MB by default (dfs.blocksize). The last block of a file is only as large as the data in it, so a 1 MB file uses 1 MB of disk, not 128 MB. Large blocks keep the metadata small and make sequential reads efficient; the price is that each block is a unit of parallelism, so a tiny file gives a job one task with almost nothing to do.
NameNode. One active process holds the entire namespace in heap: the directory tree, file attributes, and the mapping from each file to its blocks. It does not persist which DataNodes hold each block; it rebuilds that map from block reports every time it starts, which is why a restarting NameNode sits in safe mode until enough DataNodes have reported. Every namespace change is appended to an edit log before it is acknowledged, and the edit log is periodically merged into a checkpoint image (fsimage).
DataNodes. One per worker. A DataNode stores block replicas as ordinary files on local disks, each with a checksum file, and sends a heartbeat every 3 seconds (dfs.heartbeat.interval). Heartbeat replies carry commands: replicate this block, delete that one. With default settings a DataNode is declared dead after 10 minutes 30 seconds without a heartbeat (twice the 5-minute recheck interval plus ten heartbeats). The delay is deliberate: re-replicating a dead node's 50 TB is expensive, and a node that is merely rebooting should not trigger it.
Replication and placement. The default factor is 3 (dfs.replication). With rack awareness configured, the first replica goes on the writer's node (or a random node if the client is outside the cluster), the second on a node in a different rack, and the third on a different node in that second rack. The result survives the loss of a whole rack while sending only one copy across the core switch on every write.
HDFS: blocks, NameNode and DataNodes
Blocks. HDFS stores every file as a sequence of fixed-size blocks, 128 MB by default (dfs.blocksize). The last block of a file is only as large as the data in it, so a 1 MB file uses 1 MB of disk, not 128 MB. Large blocks keep the metadata small and make sequential reads efficient; the price is that each block is a unit of parallelism, so a tiny file gives a job one task with almost nothing to do.
NameNode. One active process holds the entire namespace in heap: the directory tree, file attributes, and the mapping from each file to its blocks. It does not persist which DataNodes hold each block; it rebuilds that map from block reports every time it starts, which is why a restarting NameNode sits in safe mode until enough DataNodes have reported. Every namespace change is appended to an edit log before it is acknowledged, and the edit log is periodically merged into a checkpoint image (fsimage).
DataNodes. One per worker. A DataNode stores block replicas as ordinary files on local disks, each with a checksum file, and sends a heartbeat every 3 seconds (dfs.heartbeat.interval). Heartbeat replies carry commands: replicate this block, delete that one. With default settings a DataNode is declared dead after 10 minutes 30 seconds without a heartbeat (twice the 5-minute recheck interval plus ten heartbeats). The delay is deliberate: re-replicating a dead node's 50 TB is expensive, and a node that is merely rebooting should not trigger it.
Replication and placement. The default factor is 3 (dfs.replication). With rack awareness configured, the first replica goes on the writer's node (or a random node if the client is outside the cluster), the second on a node in a different rack, and the third on a different node in that second rack. The result survives the loss of a whole rack while sending only one copy across the core switch on every write.
How bytes move: the write and read paths
Write path. The client asks the NameNode to create the file, then to allocate a block. The NameNode returns a block ID and an ordered list of three DataNodes. The client streams data in packets (64 KB by default, each made of 512-byte checksummed chunks) to the first DataNode, which forwards to the second, which forwards to the third: a pipeline, not three separate uploads. Acknowledgements flow back up the pipeline, and the client keeps a queue of unacknowledged packets so it can rebuild the pipeline without the failed node if one dies mid-write. When the file is closed, the NameNode records it complete once the minimum number of replicas has reported the final block.
Read path. The client asks the NameNode for block locations, sorted by network distance from the reader, then reads directly from the closest DataNode, verifying checksums as it goes. A checksum mismatch makes the client try another replica and report the corrupt one, which the NameNode then re-replicates from a good copy. When the reader is on the same machine as the replica, short-circuit local reads let it open the block file directly through a Unix domain socket hand-off, bypassing the DataNode's TCP path.
The design principle in both paths is the same: the NameNode handles metadata only. Bytes never pass through it, which is what lets one metadata process serve thousands of DataNodes.
NameNode high availability
A single NameNode is a single point of failure, so production clusters run two (or more, since Hadoop 3) NameNodes with a Quorum Journal Manager. The active NameNode writes every edit to a set of JournalNodes, usually three, and an edit counts as durable once a majority acknowledge it. The standby tails the journal, applies edits to its own in-memory namespace, and produces the periodic checkpoint, which is why a separate Secondary NameNode is not used in an HA setup.
Failover is driven by a ZooKeeper Failover Controller (ZKFC) beside each NameNode. The ZKFC health-checks its local NameNode and holds a ZooKeeper session; if the active's session expires, the other ZKFC wins the election and promotes its standby. Split-brain is prevented at the journal: each new active obtains a higher epoch number, and JournalNodes reject writes from any lower epoch, so a paused old active that wakes up cannot corrupt the log. DataNodes heartbeat to all NameNodes so the standby already has a fresh block map when it takes over.
<!-- hdfs-site.xml: the HA keys every cluster needs -->
<property><name>dfs.nameservices</name><value>prod</value></property>
<property><name>dfs.ha.namenodes.prod</name><value>nn1,nn2</value></property>
<property><name>dfs.namenode.shared.edits.dir</name>
<value>qjournal://jn1:8485;jn2:8485;jn3:8485/prod</value></property>
<property><name>dfs.ha.automatic-failover.enabled</name><value>true</value></property>
<property><name>dfs.client.failover.proxy.provider.prod</name>
<value>org.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider</value></property>Clients address the logical name hdfs://prod and the proxy provider retries against whichever NameNode is active. Check the state with hdfs haadmin -getServiceState nn1.
YARN: ResourceManager, NodeManagers and ApplicationMasters
YARN separates cluster-wide resource arbitration from per-application coordination.
ResourceManager. Holds the cluster's free capacity and a pluggable scheduler (Capacity Scheduler or Fair Scheduler) that decides which queue and application get the next container. It does not run or monitor tasks. With HA, an active and standby RM elect a leader through ZooKeeper and keep application state in a ZooKeeper-backed state store, so running applications survive an RM failover.
NodeManager. One per worker. It advertises how much it can give out (yarn.nodemanager.resource.memory-mb and yarn.nodemanager.resource.cpu-vcores), launches containers when given a launch context, monitors their memory, and serves shuffle data through auxiliary services such as the MapReduce or Spark external shuffle service.
Container. Not a Docker image by default: a container is simply a process tree with a resource grant. The NodeManager can enforce that grant by polling memory usage and killing offenders, by placing the process in Linux cgroups through the LinuxContainerExecutor, or by running it under the Docker runtime when that is configured.
ApplicationMaster. Each application gets its own AM, which runs in the first container granted. The AM negotiates further containers, tracks its tasks and retries failed ones. Putting this logic in a per-application process is what lets YARN host frameworks it knows nothing about.
One job from submit to commit
Follow a Spark job that reads a 10 GB file, which HDFS stores as 80 blocks.
- The client uploads the application's jars and configuration to an HDFS staging directory and submits an application to the RM, naming a queue.
- The scheduler grants a container for the AM on some NodeManager; the NM localises the staged files and starts the driver or AM process.
- The AM asks the NameNode for the 80 blocks' locations and turns them into resource requests that name preferred nodes, preferred racks and
*(anywhere). - On each NodeManager heartbeat the scheduler tries to match a pending request to that node. With the Capacity Scheduler, delay scheduling skips up to
yarn.scheduler.capacity.node-locality-delayscheduling opportunities (default 40) waiting for a node-local slot before relaxing to rack-local. - Tasks read their blocks, mostly from local disk. Shuffle output is served by the NodeManagers.
- Output is written to a temporary directory and committed by renames, the AM unregisters, and the NodeManagers aggregate container logs into HDFS if log aggregation is on.
If the AM container dies, the RM restarts it up to yarn.resourcemanager.am.max-attempts times (default 2), so frameworks that cannot recover their progress simply rerun.
Worked example: sizing a 40-node cluster
Size a 40-node cluster where each worker has 12 x 8 TB disks, 256 GB RAM and 48 hardware threads.
Storage. Raw capacity is 40 x 96 TB = 3.84 PB. At replication 3 that is 1.28 PB of logical data; keep 25 percent free for re-replication after failures and for balancer moves, leaving about 960 TB usable.
NameNode heap. A widely used rule of thumb is about 150 bytes of heap per namespace object (file, directory or block), and roughly 1 GB of heap per million blocks once overheads are included. If the average file is 64 MB, 960 TB is about 15 million files, each one block, so about 30 million objects and 15 million blocks: a 16 GB heap fits, and 32 GB gives room to grow and to collect garbage calmly. If the same data arrives as 1 MB files, that becomes nearly a billion files and close to two billion objects, which no single NameNode heap holds. The small-files problem is a metadata problem first.
YARN memory. Reserve memory for the OS, DataNode and NodeManager, say 32 GB, and give YARN the rest: yarn.nodemanager.resource.memory-mb=229376 (224 GB). With 4 GB containers that is 56 containers per node by memory, but only 44 vcores remain after reserving four, so with the DominantResourceCalculator cores become the limit at 44. The Capacity Scheduler's default calculator counts memory only, which is how clusters end up running 56 CPU-bound tasks on 44 threads.
def namenode_heap_gb(logical_tb, avg_file_mb, blocks_per_file=1.0):
files = logical_tb * 1024 * 1024 / avg_file_mb
blocks = files * blocks_per_file
return files, blocks / 1e6 # ~1 GB heap per million blocks
def containers_per_node(yarn_mem_gb, yarn_vcores, mem_gb, vcores):
return min(yarn_mem_gb // mem_gb, yarn_vcores // vcores)
print(namenode_heap_gb(960, 64)) # ~15.7M files -> ~16 GB heap
print(namenode_heap_gb(960, 1)) # ~1.0B files -> not one NameNode
print(containers_per_node(224, 44, 4, 1)) # 44, cores bind first
Failure modes
These are the incidents that recur across Hadoop shops, roughly in order of frequency.
- NameNode garbage-collection pauses. A long pause misses the ZKFC health check, triggers a failover, and the cluster flaps. Size the heap from the arithmetic above, use a low-pause collector, and alert on GC time before it reaches the ZooKeeper session timeout.
- Re-replication storms. Losing a rack leaves millions of under-replicated blocks. The NameNode throttles the work per DataNode, so recovery takes hours, and during it a second failure can lose data. Watch
UnderReplicatedBlocksandMissingBlocks. - One bad disk takes a DataNode down.
dfs.datanode.failed.volumes.tolerateddefaults to 0, so on dense nodes set it to 1 or 2 rather than losing 96 TB of replicas to one drive. - Small files. Heap pressure, slow block reports and thousands of tiny tasks. Compact into larger files or table formats.
- Unhealthy NodeManagers. Full local directories or a failing health script mark the node unhealthy and YARN stops scheduling there, quietly shrinking the cluster.
- Queue starvation. One queue at its maximum capacity while others sit idle, usually a configuration with no elasticity or no preemption.
Operating the cluster
A small set of commands answers most operational questions:
hdfs dfsadmin -report # capacity, live and dead DataNodes
hdfs fsck / -list-corruptfileblocks # what is actually lost
hdfs haadmin -getAllServiceState # which NameNode is active
hdfs balancer -threshold 10 # even out disk use after adding nodes
yarn node -list -all # NodeManager states, including UNHEALTHY
yarn application -list -appStates RUNNING
yarn rmadmin -refreshQueues # apply capacity-scheduler.xml changesAlert on missing blocks (page immediately), under-replicated blocks (trend), NameNode RPC queue time, NameNode GC time, dead and unhealthy nodes, and pending containers per queue. Decommission nodes through the exclude file rather than switching them off, so their replicas are copied away first, and run the balancer after expansions because HDFS does not move existing data on its own.
Trade-offs
| Decision | Option A | Option B |
|---|---|---|
| Durability | 3x replication: fast reads, locality, 200% overhead | Erasure coding RS-6-3: 50% overhead, survives 3 losses, no locality and costly rebuilds |
| Compute placement | Co-located with HDFS: data locality, coupled scaling | Disaggregated over object storage: independent scaling, network-bound reads |
| Scheduler | Capacity: guaranteed queue shares, predictable | Fair: shares converge dynamically, better for ad hoc mixes |
| Namespace scale | One HA nameservice: simple | Federation or router-based federation: more objects, more moving parts |
Use erasure coding for cold data and keep replication for hot data that jobs read repeatedly. Choose co-location when your workloads are scan-heavy and steady; choose object storage when compute demand is bursty and storage grows faster than compute.
What to do next
To go deeper, read the NameNode internals, NameNode high availability, rack awareness, the ResourceManager, the Capacity Scheduler and HDFS erasure coding. Then work through this checklist:
- Confirm HA: two NameNodes, three JournalNodes, automatic failover on, and a tested manual failover.
- Configure the rack topology script and verify placement with
hdfs fsck <path> -files -blocks -racks. - Compute NameNode heap from your block count and alert on GC time and RPC queue time.
- Set YARN memory and vcores per node, and switch to the DominantResourceCalculator if tasks are CPU-bound.
- Find directories with a low average file size and schedule compaction.
- Set a failed-volume tolerance on dense nodes and rehearse decommissioning one node.