A DataNode is the HDFS process that owns disks. It stores block replicas as ordinary files, serves reads and pipelined writes, verifies checksums in the background, and keeps the NameNode's picture of where every replica lives accurate. It knows nothing about file names or directories; that namespace lives in the NameNode. Most HDFS incidents that are not NameNode incidents are DataNode incidents: a failing disk, a slow node dragging a write pipeline, a block report storm after a restart, or a decommission that never finishes.
This article walks through the DataNode from the disk up: on-disk layout, replica states, the heartbeat and block report protocol, the data path, integrity checking, and the operational decisions that follow. Configuration names and defaults are for Hadoop 3; verify them against your distribution. An earlier version of this page said full block reports are sent about hourly. The Hadoop 3 default for dfs.blockreport.intervalMsec is 21,600,000 ms, which is six hours.
On-disk layout
Each entry in dfs.datanode.data.dir is a volume, normally one physical disk, optionally tagged with a storage type such as [DISK], [SSD] or [ARCHIVE]. Inside each volume, a current/VERSION file records the storage ID and cluster ID, and each block pool (one per NameNode namespace in a federated cluster) gets its own directory:
/data/3/dfs/dn/current/
VERSION # storageID, clusterID, layout version
BP-1832746105-10.0.4.11-1690000000000/
current/
finalized/subdir0/subdir17/
blk_1073742913 # the block bytes, a plain file
blk_1073742913_2089.meta # header + one CRC per 512-byte chunk
rbw/ # replicas being written
tmp/ # replicas being copied inThe number after the block ID in the meta file name is the generation stamp, which the NameNode bumps whenever a replica is recovered or appended so that stale copies can be detected. The meta file holds one checksum per dfs.bytes-per-checksum bytes (512 by default), CRC32C unless configured otherwise. Four bytes per 512 is about 0.8 percent overhead, and it is what lets the DataNode and client verify every chunk they move. Blocks are spread over nested subdirectories so no directory grows huge. Never edit these trees by hand; the DataNode's in-memory replica map is the source of truth while it runs.
On-disk layout
Each entry in dfs.datanode.data.dir is a volume, normally one physical disk, optionally tagged with a storage type such as [DISK], [SSD] or [ARCHIVE]. Inside each volume, a current/VERSION file records the storage ID and cluster ID, and each block pool (one per NameNode namespace in a federated cluster) gets its own directory:
/data/3/dfs/dn/current/
VERSION # storageID, clusterID, layout version
BP-1832746105-10.0.4.11-1690000000000/
current/
finalized/subdir0/subdir17/
blk_1073742913 # the block bytes, a plain file
blk_1073742913_2089.meta # header + one CRC per 512-byte chunk
rbw/ # replicas being written
tmp/ # replicas being copied inThe number after the block ID in the meta file name is the generation stamp, which the NameNode bumps whenever a replica is recovered or appended so that stale copies can be detected. The meta file holds one checksum per dfs.bytes-per-checksum bytes (512 by default), CRC32C unless configured otherwise. Four bytes per 512 is about 0.8 percent overhead, and it is what lets the DataNode and client verify every chunk they move. Blocks are spread over nested subdirectories so no directory grows huge. Never edit these trees by hand; the DataNode's in-memory replica map is the source of truth while it runs.
Replica states
A replica moves through states that explain most of what you see in logs. TEMPORARY replicas are being copied for re-replication or balancing and are invisible to readers. RBW, replica being written, belongs to an open file in a client write pipeline; it lives in rbw/ and its visible length grows as packets are acknowledged. FINALIZED replicas are complete and immutable until an append. RWR, replica waiting to be recovered, is what an RBW replica becomes when the DataNode restarts mid-write. RUR, replica under recovery, marks a replica taking part in lease or block recovery, where the primary DataNode agrees a final length and generation stamp with its peers. The detailed recovery protocol is covered in the HDFS write pipeline.
Heartbeats, block reports and liveness
Every DataNode-to-NameNode conversation is started by the DataNode. One BPServiceActor thread per NameNode per block pool sends a heartbeat every dfs.heartbeat.interval (3 seconds) with capacity, usage, failed volumes and active transfer count. The NameNode never calls the DataNode; it answers heartbeats with commands: replicate this block to that node, delete these blocks, recover this block, cache or uncache. That design keeps the NameNode from ever blocking on a slow DataNode.
Block state reaches the NameNode in two forms. Incremental block reports go out shortly after a replica is received or deleted, so new data becomes readable quickly. Full block reports list every replica on every volume every six hours by default, and let the NameNode reconcile drift. When a node holds more than dfs.blockreport.split.threshold blocks (1,000,000 by default), the report is sent one storage at a time to keep each RPC bounded.
Liveness has two thresholds. After dfs.namenode.stale.datanode.interval (30 seconds) without a heartbeat a node is marked stale, and can be deprioritized for reads and writes if the avoid-stale settings are enabled. A node is declared dead after 2 times dfs.namenode.heartbeat.recheck-interval plus 10 times the heartbeat interval: 2 x 300 s + 10 x 3 s = 630 seconds, ten and a half minutes. Only then does re-replication of its blocks begin. That delay is deliberate: a rebooting node should not trigger the copy of tens of terabytes.
The data path
Data moves over the DataNode's streaming protocol on dfs.datanode.address (port 9866 in Hadoop 3), not over RPC. Each connection is served by a DataXceiver thread, bounded by dfs.datanode.max.transfer.threads (4,096 by default). Writes arrive as packets of about 64 KB containing 512-byte chunks and their checksums; the DataNode verifies, writes to disk, forwards to the next node, and relays acknowledgements back. Reads stream the block file and its checksums; the client verifies and reports corrupt replicas to the NameNode.
When the client runs on the same host, short-circuit local reads let it bypass the socket: the DataNode passes open file descriptors over a Unix domain socket at dfs.domain.socket.path, and the client reads the file directly. For colocated compute engines this is one of the largest read-performance settings available. Its tuning, along with hedged reads, is covered in HDFS performance tuning.
Integrity: scanners and disk checks
Disks fail gradually, so a DataNode runs three kinds of integrity check. A VolumeScanner per volume reads every block and checks its checksums over dfs.datanode.scan.period.hours (504 hours, three weeks), throttled by dfs.block.scanner.volume.bytes.per.second (1 MB/s by default). Blocks that clients recently reported suspect are scanned first. Corrupt replicas are reported to the NameNode, which re-replicates from a good copy and then deletes the bad one. On a 16 TB disk the default throttle cannot finish a full pass in three weeks, so dense nodes should raise it or accept a longer effective period knowingly.
The DirectoryScanner runs every dfs.datanode.directoryscan.interval (21,600 seconds) and compares the in-memory replica map with what is actually on disk, fixing missing files and orphans. The disk checker runs when I/O errors occur and marks a volume failed if it cannot be read or written. What happens next is set by dfs.datanode.failed.volumes.tolerated, which defaults to 0: the first failed disk shuts the whole DataNode down. On a node with twelve disks that turns one dead disk into a whole-node re-replication, so most operators raise it.
A starting configuration
A reasonable starting point for a dense, twelve-disk node:
<property><name>dfs.datanode.data.dir</name>
<value>[DISK]/data/1/dfs/dn,[DISK]/data/2/dfs/dn,...,[DISK]/data/12/dfs/dn</value></property>
<property><name>dfs.datanode.failed.volumes.tolerated</name><value>2</value></property>
<property><name>dfs.datanode.du.reserved</name><value>107374182400</value></property> <!-- 100 GB per volume -->
<property><name>dfs.datanode.fsdataset.volume.choosing.policy</name>
<value>org.apache.hadoop.hdfs.server.datanode.fsdataset.AvailableSpaceVolumeChoosingPolicy</value></property>
<property><name>dfs.datanode.max.transfer.threads</name><value>8192</value></property>
<property><name>dfs.client.read.shortcircuit</name><value>true</value></property>
<property><name>dfs.domain.socket.path</name><value>/var/lib/hadoop-hdfs/dn_socket</value></property>The default round-robin volume policy ignores free space, so a disk replaced half way through a node's life stays emptier than its neighbours; the available-space policy biases new blocks toward emptier volumes. For existing skew inside a node, use the disk balancer; for skew between nodes, the cluster balancer. Reserved space keeps non-HDFS writers such as logs and shuffle data from filling a disk to zero.
Worked example: one disk fails on a dense node
A 200-node cluster stores data with three-way replication. Each node has twelve 16 TB disks at 70 percent full, about 11.2 TB of replicas per disk. At 02:00, disk 7 on node dn117 starts returning I/O errors.
With the default failed.volumes.tolerated = 0, the DataNode shuts down. Ten and a half minutes later the NameNode declares it dead, and about 134 TB of replicas, every block on all twelve disks, become under-replicated. The NameNode schedules copies from surviving replicas across the cluster, limited per source node by its replication-stream settings. Spread over 199 nodes that is about 0.68 TB of reads and writes each, hours of background I/O competing with production jobs, all for one disk.
With failed.volumes.tolerated = 2, the DataNode marks disk 7 failed, keeps serving the other eleven, and reports the loss in its next heartbeat. Only about 11.2 TB becomes under-replicated: twelve times less re-replication. The on-call engineer sees NumFailedVolumes rise, replaces the disk, adds it back to dfs.datanode.data.dir, and hot-swaps it without a restart:
hdfs dfsadmin -reconfig datanode dn117.example.com:9867 start
hdfs dfsadmin -reconfig datanode dn117.example.com:9867 status
hdfs dfsadmin -triggerBlockReport dn117.example.com:9867
hdfs diskbalancer -plan dn117.example.com # then -execute the generated plan fileThe disk balancer step matters: the new disk is empty while its neighbours are 70 percent full, and without it reads concentrate on the old disks for months.
Monitoring and planned maintenance
Watch DataNodes through the NameNode's view and their own JMX. hdfs dfsadmin -report shows per-node capacity, usage and last contact; hdfs fsck / -list-corruptfileblocks shows damage. Each DataNode serves JMX on its HTTP port (9864):
import json, urllib.request
def dn_health(host, port=9864):
with urllib.request.urlopen(f"http://{host}:{port}/jmx", timeout=10) as r:
beans = json.load(r)["beans"]
out = {}
for b in beans: # bean and attribute names vary by version: list /jmx once and confirm
name = b.get("name", "")
if name.startswith("Hadoop:service=DataNode,name=FSDatasetState"):
out["failed_volumes"] = b.get("NumFailedVolumes")
out["remaining_bytes"] = b.get("Remaining")
elif name.startswith("Hadoop:service=DataNode,name=DataNodeActivity"):
out["write_block_avg_ms"] = b.get("WriteBlockOpAvgTime")
out["read_block_avg_ms"] = b.get("ReadBlockOpAvgTime")
out["heartbeat_avg_ms"] = b.get("HeartbeatsAvgTime")
return outAlert on failed volumes, on nodes whose heartbeat age passes 30 seconds, on DataXceiver counts near the transfer-thread limit, and on one node whose write-block time sits far above its peers, because a single slow disk slows every pipeline that includes it. For planned work, prefer maintenance state over decommissioning when a node will be back within hours: decommissioning copies every replica elsewhere, while maintenance only keeps a minimum replica count and avoids the copy.
Failure modes
- One bad disk takes the node down because failed volumes tolerated is still 0.
- Block report storms after a rack restart, when many nodes send full reports together and NameNode RPC queues back up. Stagger restarts.
- Xceiver exhaustion from many small concurrent reads, causing failed reads elsewhere.
- Silent bit rot on cold data that the throttled scanner has not reached yet; clients catch it on read, but only for data that is actually read.
- Stuck decommissioning when replicas of files still open for write, or a cluster short of free space, keep the last blocks from moving. Check
hdfs dfsadmin -report -decommissioningand list open files before waiting longer. - Slow-disk tail latency that never fails health checks but dominates p99 of pipelines.
- Disk full because reserved space was not set and non-HDFS data grew on shared volumes.
Trade-offs
HDFS chooses JBOD over RAID on purpose: replication across nodes already protects data, so RAID would spend disks twice and lose the per-disk parallelism. Denser nodes are cheaper per terabyte but turn each node failure into a larger re-replication and longer windows of reduced durability. Erasure coding cuts storage overhead from 200 to 50 percent for cold data, at the cost of reconstruction reads that touch several DataNodes. Short-circuit reads speed colocated compute but tie the client to local file access and its permissions.
What to do next
- Check
dfs.datanode.failed.volumes.toleratedon every node and raise it on multi-disk nodes. - Compute whether the block scanner throttle can cover your disk size within the scan period.
- Set reserved space per volume and switch to the available-space volume policy if disks were replaced.
- Export failed volumes, heartbeat age, xceiver count and per-node write latency from JMX, and alert on them.
- Practise a disk hot-swap with
-reconfigand a disk balancer run on a staging node. - Use maintenance state, not decommission, for short planned outages.