HDFS durability comes from a simple rule enforced continuously: every block of every file should have a target number of replicas, spread across failure domains, and the NameNode keeps repairing the cluster until reality matches that target. The write path places the first replicas, but most of the work happens afterwards, when disks fail, nodes are retired, racks lose power and corrupted bytes turn up during a read. This article follows a block through its whole life: what it is on disk, where its replicas go, how the NameNode learns what exists, how it prioritises repairs, and what the configuration knobs actually do. Every property name and default below was checked against the current stable hdfs-default.xml.

The write pipeline itself, with its packet acknowledgements and recovery, is covered in the HDFS write pipeline article, and topology scripts in rack awareness. Here we focus on the replica lifecycle that sits on top of both.

One block, three replicas, two racksNameNodeblocks map + redundancy monitorRack ADN a1 (writer)replica 1: blk_1073741825DN a2other blocksDN a3other blocksRack BDN b1replica 2DN b2replica 3DN b3Rack CDN c1re-replication targetDN c2pipelinecopy if b1 diesheartbeats, block reportsreplicate commandsEach replica = blk_ID data file + blk_ID_genstamp.meta (CRC32C per 512 bytes)Rack rule: no block may live on a single rack while two racks existMonitor rule: under-replicated blocks are queued by how close they are to loss
Replicas span two racks; the NameNode learns state from heartbeats and block reports and schedules repairs onto a third rack when a holder fails.

What a block is, on disk and in the NameNode

A file in HDFS is a sequence of blocks. The block size is a per-file attribute fixed at creation, defaulting to dfs.blocksize = 134217728 bytes (128 MiB). Every block except the last is full; the last holds whatever remains, so a 300 MiB file is two 128 MiB blocks plus a 44 MiB block, and the 44 MiB block uses 44 MiB of disk per replica, not 128. Block size is a metadata decision, not a space reservation.

Each block has a 64-bit ID and a generation stamp. The generation stamp increments whenever the block is recovered or reopened for append, which lets the NameNode tell a current replica from a stale one that missed an update while its DataNode was unreachable. A replica on a DataNode is two ordinary files in a local directory: blk_1073741825 holds the bytes and blk_1073741825_1001.meta holds a header plus one checksum per dfs.bytes-per-checksum (512) bytes, computed with dfs.checksum.type (CRC32C by default). Clients verify these checksums on every read.

Distinguish block from replica. The NameNode stores the block (ID, length, generation stamp, owning file, expected replication). DataNodes store replicas. The NameNode's view of which replicas exist is not persisted in the fsimage at all: it is rebuilt in memory from DataNode reports after every restart, which is why a freshly started NameNode sits in safe mode until enough blocks have been reported.

Why 128 MiB? Fewer, bigger blocks mean fewer NameNode objects (a commonly cited rule of thumb is about 150 bytes of heap per file, directory or block object, so treat it as an estimate) and long sequential reads. Smaller blocks mean more parallelism for jobs that split work by block. Most clusters keep the default and raise it per file only for very large, scan-heavy datasets.

What a block is, on disk and in the NameNode

A file in HDFS is a sequence of blocks. The block size is a per-file attribute fixed at creation, defaulting to dfs.blocksize = 134217728 bytes (128 MiB). Every block except the last is full; the last holds whatever remains, so a 300 MiB file is two 128 MiB blocks plus a 44 MiB block, and the 44 MiB block uses 44 MiB of disk per replica, not 128. Block size is a metadata decision, not a space reservation.

Each block has a 64-bit ID and a generation stamp. The generation stamp increments whenever the block is recovered or reopened for append, which lets the NameNode tell a current replica from a stale one that missed an update while its DataNode was unreachable. A replica on a DataNode is two ordinary files in a local directory: blk_1073741825 holds the bytes and blk_1073741825_1001.meta holds a header plus one checksum per dfs.bytes-per-checksum (512) bytes, computed with dfs.checksum.type (CRC32C by default). Clients verify these checksums on every read.

Distinguish block from replica. The NameNode stores the block (ID, length, generation stamp, owning file, expected replication). DataNodes store replicas. The NameNode's view of which replicas exist is not persisted in the fsimage at all: it is rebuilt in memory from DataNode reports after every restart, which is why a freshly started NameNode sits in safe mode until enough blocks have been reported.

Why 128 MiB? Fewer, bigger blocks mean fewer NameNode objects (a commonly cited rule of thumb is about 150 bytes of heap per file, directory or block object, so treat it as an estimate) and long sequential reads. Smaller blocks mean more parallelism for jobs that split work by block. Most clusters keep the default and raise it per file only for very large, scan-heavy datasets.

Where replicas go

The default placement policy (BlockPlacementPolicyDefault) chooses targets when a block is allocated and again when a repair needs a new home:

  1. First replica on the writer's own DataNode if the client runs on one; otherwise on a random node, skipping nodes that are full, decommissioning or too busy.
  2. Second replica on a node in a different rack.
  3. Third replica on a different node in the same rack as the second.
  4. Any further replicas on random nodes, with a cap of (replicas - 1) / racks + 2 replicas per rack.

This layout is a deliberate compromise. Two racks hold the block, so losing a whole rack never loses it. Only one hop of the write pipeline crosses the rack boundary, so write traffic on the core network is one copy rather than two. Reads find a replica on the writer's rack for free. The cost is that two of the three replicas share a rack, so a rack failure leaves those blocks with a single surviving copy, a point the worked example below makes concrete.

Two more inputs shape the choice. With dfs.namenode.redundancy.considerLoad = true (the default) the NameNode avoids nodes whose active transfer count is well above the cluster average. And storage policies (HOT, WARM, COLD, ONE_SSD, ALL_SSD and others) restrict which storage types on a node may hold each replica, so placement chooses node and storage type together.

How the NameNode learns what exists

The NameNode never polls disks. It learns state from three DataNode messages:

  • Heartbeats every dfs.heartbeat.interval = 3 seconds, carrying capacity, usage and transfer counts. The reply carries commands: replicate this block to that node, delete these replicas.
  • Incremental block reports, sent soon after a replica is received or deleted, so the NameNode learns about new replicas within seconds.
  • Full block reports every dfs.blockreport.intervalMsec = 21600000 ms (six hours), listing every replica, which repairs any drift between the two views.

Failure detection is time-based. A node that misses heartbeats for dfs.namenode.stale.datanode.interval = 30 seconds is marked stale. It is declared dead after 2 x dfs.namenode.heartbeat.recheck-interval + 10 x heartbeat interval = 2 x 300 s + 10 x 3 s = 630 seconds, ten and a half minutes. Only then do its replicas stop counting. The long window is intentional: a node rebooting for a kernel patch should not trigger terabytes of copying.

Stale state is only useful if you act on it. dfs.namenode.avoid.read.stale.datanode and dfs.namenode.avoid.write.stale.datanode both default to false; turning them on moves stale nodes to the end of read location lists and out of new write pipelines during those ten minutes, which removes a familiar class of slow reads and stuck writes.

The redundancy monitor and its priority queues

Repairs are driven by the redundancy monitor, a NameNode thread that wakes every dfs.namenode.redundancy.interval.seconds = 3 seconds. Blocks needing work live in a set of priority queues, ordered by how close each block is to being lost:

PriorityConditionMeaning
0, highestone live replica left, or only decommissioning replicasone more failure loses data
1, very low redundancylive replicas below one third of the targetbadly degraded
2, low redundancybelow target, otherwiseordinary repair
3, badly distributedenough replicas, but all on one rackrack rule broken
4, corruptevery replica is corruptneeds an operator

Each pass takes up to dfs.namenode.replication.work.multiplier.per.iteration (2) x live DataNodes blocks off the queues, highest priority first, picks a source replica and a target, and hands the work out via heartbeat replies. Per-source concurrency is capped by dfs.namenode.replication.max-streams = 2 for normal work, and by the hard limit of 4 for highest-priority blocks. Scheduled work that has not been confirmed by an incremental block report within dfs.namenode.reconstruction.pending.timeout-sec = 300 seconds goes back on the queue. The loop, simplified:

every redundancy_interval (3 s):
    budget = work_multiplier * live_datanodes
    for block in low_redundancy_queues.ordered_by_priority():
        if budget == 0: break
        src = choose_source(block)        # live, not corrupt, under its stream limit;
                                          # decommissioning holders may serve as sources
        if src is None: continue          # every holder busy: retry next pass
        targets = placement.choose(block, needed = expected - live - pending)
        src.pending_commands.add(Replicate(block, targets))
        pending[block] = now()
        budget -= 1
    for block, t in pending.items():
        if now() - t > pending_timeout:   # 300 s, no block report arrived
            requeue(block)

Worked example: losing a node, then a rack

Take a 200-node cluster in 10 racks of 20. Each DataNode holds about 28 TB of block data, or roughly 215,000 replicas at 128 MiB each. Node a7 loses its motherboard at 10:00:00.

  1. 10:00:30: a7 is stale. With the avoid-stale settings on, readers skip it and new pipelines exclude it.
  2. 10:10:30: a7 is dead. Its 215,000 blocks now show 2 of 3 live replicas. Two thirds is not below one third, so they go to priority 2.
  3. Each 3-second pass may schedule 2 x 199 = 398 blocks, about 133 blocks per second or roughly 17 GB/s of copying. That is a scheduling ceiling, not a forecast: at that rate the backlog would clear in about 27 minutes, but each source sends only two streams at a time and real disks and links usually stretch recovery to somewhere between half an hour and a few hours.
  4. Because a7's blocks have their other replicas scattered over every other rack, the copying is spread over the whole cluster instead of hammering one neighbour. This many-to-many repair is the main reason HDFS recovers from node loss faster than a mirrored pair would.

Now lose all of rack A at once: 20 nodes, about 4.3 million replicas. No block is lost, because every block spans two racks. But consider blocks whose second and third replicas lived on rack A, because a writer on another rack chose A as its remote rack: they are now down to one replica and land in priority 0. The monitor repairs those first, which is exactly what you want, and the max-streams-hard-limit exists so that this urgent work can exceed the normal per-node cap. Watch the priority 0 count in a rack drill; it is the number that tells you how exposed you were.

Over-replication, corruption, missing blocks and decommissioning

Over-replication. When a dead node returns with its disks intact, or after hdfs dfs -setrep lowers a file's target, blocks have excess replicas. The NameNode picks which to delete while preserving rack spread: it removes from a rack holding more than one replica, preferring nodes with less free space, and sends delete commands in heartbeat replies.

Corruption. A checksum mismatch found by a reader, or by the DataNode block scanner (each replica is re-verified roughly every dfs.datanode.scan.period.hours = 504 hours, three weeks), is reported to the NameNode. The corrupt replica stops counting as live, a good replica is copied, and only then is the bad one deleted, so the count never drops further than it must. If every replica is corrupt the block goes to priority 4 and the file is listed by hdfs fsck / -list-corruptfileblocks.

Missing blocks. A block with no live replicas at all is missing. Reads fail. Often the replicas are on nodes that are down rather than destroyed, so bring nodes back before deleting files.

Decommissioning. Adding a host to the exclude file and running hdfs dfsadmin -refreshNodes moves it to Decommission In Progress. Its replicas stop counting toward the target, so the monitor copies every block elsewhere, using the retiring node as a source where possible. It reaches Decommissioned only when all of its blocks are fully replicated elsewhere. Maintenance state is the short-outage alternative: the node may go down while each of its blocks keeps at least dfs.namenode.maintenance.replication.min (1) live replica elsewhere, avoiding a full copy for a planned reboot.

Operating it: commands and metrics

# Where are the replicas of one file, and are any missing?
hdfs fsck /data/events/2026-10-03 -files -blocks -locations -racks

# Block size (%o) and replication (%r) of a file
hdfs dfs -stat "%o %r %n" /data/events/2026-10-03/part-00000.parquet

# Cluster-wide view: dead and decommissioning nodes, low-redundancy counts
hdfs dfsadmin -report -dead
hdfs dfsadmin -metasave meta.txt     # dumps queues to the NameNode log directory

# Lower replication of cold data and wait for the excess to be removed
hdfs dfs -setrep -w 2 /data/archive/2024

# Files with corrupt blocks
hdfs fsck / -list-corruptfileblocks

For alerting, the NameNode exports MissingBlocks, CorruptBlocks, UnderReplicatedBlocks (LowRedundancyBlocks in newer releases) and PendingReplicationBlocks. Page on any missing block. Alert when low-redundancy counts stop falling after a failure, because a flat line means the monitor cannot find sources or targets: full disks, a rack that is too small to satisfy the rack rule, or every holder at its stream limit.

Trade-offs

ChoiceGainCost
Replication 3 (default)survives any two node losses or one rack loss; fast local reads200% storage overhead
Replication 2one third less diska single failure leaves one copy; risky during the 10.5 minute detection window
Replication above 3 for hot filesmore read locality for widely shared datadisk, and longer pipelines on write
Erasure coding RS-6-350% overhead, tolerates three lost cellsCPU for encoding, reconstruction reads from six nodes, no locality
Bigger blocksless NameNode heap, longer sequential I/Ofewer splits, more data per repair unit

The usual pattern is replication for hot and recently written data and erasure coding for cold directories. Keep the NameNode's own memory and RPC budget in mind as block counts grow; the NameNode article covers that side, and the DataNode article covers the volume and scanner details.

Failure modes

  • Two-rack clusters. With two racks, the rule that the third replica shares the second's rack still holds, but losing one rack can leave many blocks at one replica. Three racks is the practical minimum.
  • Unbalanced racks. A small rack fills first, and blocks that need it for spread queue as badly distributed. Keep racks roughly equal.
  • Replication 1 for scratch data. Fine until a scratch file becomes a dependency. Any disk failure makes it a missing block.
  • Mass restart. Restarting many DataNodes at once can push nodes past the dead threshold and start needless copying. Roll restarts or use maintenance state.
  • Deleting files with missing blocks too early. Check whether the holders are merely offline.

What to do next

  1. Run hdfs fsck / -racks and confirm the rack count is what your topology script intends.
  2. Turn on the two avoid-stale settings and watch read latency during the next node failure.
  3. Add alerts for missing and corrupt blocks, and a trend alert on low-redundancy counts.
  4. Run a rack-loss drill in staging and record the priority 0 count and time to recovery.
  5. Use maintenance state for planned reboots and decommissioning for permanent removals.
  6. Move cold directories to erasure coding and lower replication where the data is reproducible.
Key takeaway: An HDFS block is metadata on the NameNode plus checksummed replica files on DataNodes, and its durability is a control loop: placement puts replicas on two racks, heartbeats and block reports tell the NameNode what survives, and the redundancy monitor repairs the most endangered blocks first within per-node stream limits. Learn the 630-second dead-node window, enable stale-node avoidance, alert on missing and low-redundancy blocks, and use maintenance state for planned outages.