Removing a DataNode from an HDFS cluster looks trivial: stop the process. HDFS will notice and re-replicate the missing blocks. But stopping a node first means every block it held drops one replica before a copy exists elsewhere, and any block that was already at one replica becomes missing until the node comes back. Decommissioning reverses the order. The NameNode copies every block the node holds to other nodes first, while the node is still serving reads, and marks it safe to stop only when nothing depends on it.

The general procedure is covered in Hadoop node decommissioning. This article goes inside the NameNode: the admin state machine, how the monitor walks blocks, which settings throttle it, how long a decommission should take, why some get stuck at 99 percent, and when maintenance mode is the better tool.

Advertisement

Why not just stop the node

When a DataNode stops heartbeating, the NameNode waits before declaring it dead. With default settings the timeout is two recheck intervals of five minutes plus ten heartbeat intervals of three seconds, which is 10.5 minutes. During that window the node's replicas still count, so blocks look healthy while they are not. After it, every block on the node becomes under-replicated at once and the NameNode queues all of them, competing with client traffic.

Decommissioning avoids both problems. The node stays alive and serving, so replica counts never drop. The work is spread over time and governed by throttles. And the NameNode tells you, per node, how many blocks still depend on it, so the moment you stop the process is chosen from data instead of hope.

The admin state machine

Each DataNode descriptor in the NameNode carries an admin state separate from its liveness. A node is NORMAL by default. Excluding it moves it to DECOMMISSION_INPROGRESS: it stops receiving new block writes, keeps serving reads, and can act as a source for replication. When every block it holds has enough live replicas on other nodes, it becomes DECOMMISSIONED. Maintenance mode has a parallel pair, ENTERING_MAINTENANCE and IN_MAINTENANCE, covered below.

DataNode admin states as the NameNode tracks themNORMALreads + writesDECOMMISSION_INPROGRESSreads only, re-replicatingDECOMMISSIONEDsafe to stopexcludedall blocks safeENTERING_MAINTENANCEreaching min replicasIN_MAINTENANCEdown, not re-replicatedmaintenancemin reachedexpiry or removed: back to NORMALNameNode: DatanodeAdminManager monitor (every dfs.namenode.decommission.interval)scan tracked nodes' blocks -> queue low-redundancy blocks -> RedundancyMonitor schedules copiessources: any replica holder, decommissioning nodes allowed | targets: placement policy (racks, upgrade domains)complete when every block has enough live replicas elsewhere and no open file holds it back
Decommission moves a node out permanently; maintenance takes it down temporarily without re-replicating everything. Both are driven by the same NameNode monitor.

A decommissioned node that is still running keeps its data and stays visible in reports, which is useful: if something went wrong you can return it to service. Its replicas no longer count towards the replication factor once it is decommissioned.

Advertisement

Telling the NameNode

The classic mechanism is two plain files. dfs.hosts lists nodes allowed to register; dfs.hosts.exclude lists nodes to decommission. You add the hostname to the exclude file on the NameNode host, or on both NameNodes in an HA pair, and run hdfs dfsadmin -refreshNodes. Forgetting the standby NameNode is a classic mistake: after a failover the new active treats the node as normal and the decommission silently stops.

Hadoop 3 also supports a single JSON hosts file, enabled by setting dfs.namenode.hosts.provider.classname to org.apache.hadoop.hdfs.server.blockmanagement.CombinedHostFileManager and pointing dfs.hosts at the JSON file. Each entry can carry an adminState, which is how you request maintenance with an expiry time:

[
  {"hostName": "dn01.example.com"},
  {"hostName": "dn02.example.com", "adminState": "DECOMMISSIONED"},
  {"hostName": "dn03.example.com", "adminState": "IN_MAINTENANCE",
   "maintenanceExpireTimeInMS": 1790900000000}
]
hdfs dfsadmin -refreshNodes                    # run against each NameNode in an HA pair
hdfs dfsadmin -report -decommissioning         # nodes still in progress
hdfs dfsadmin -report -enteringmaintenance

The JSON form is easier to manage from configuration tooling, because the desired state of every node lives in one file instead of being inferred from two.

What the NameNode actually does

The work is split between two background loops. The admin monitor, run by DatanodeAdminManager every dfs.namenode.decommission.interval seconds (default 30), walks the blocks of each tracked node. For each block it checks whether there are enough live replicas on nodes that are not leaving. Blocks that are short are added to the NameNode's low-redundancy queues, and the node stays in progress. A node whose scan finds nothing short moves to its final state.

The second loop, the redundancy monitor, runs every dfs.namenode.redundancy.interval.seconds (default 3) and turns queued blocks into replication commands sent to DataNodes on their next heartbeat. It picks a source holding a replica, which may be the decommissioning node itself, and a target chosen by the block placement policy, so rack awareness and upgrade domains still apply. The target DataNode copies the block from the source over the data-transfer protocol and reports it, and the replica count rises. See HDFS blocks and replication for the underlying mechanics.

Walking millions of blocks holds the namespace lock, so the monitor processes about dfs.namenode.decommission.blocks.per.interval blocks (default 500,000) per run and tracks at most dfs.namenode.decommission.max.concurrent.tracked.nodes nodes (default 100) at a time; the rest wait in a pending queue. Newer releases also offer an alternative monitor, selected with dfs.namenode.decommission.monitor.class set to org.apache.hadoop.hdfs.server.blockmanagement.DatanodeAdminBackoffMonitor, which limits how many decommission blocks it loads into the replication queue (dfs.namenode.decommission.backoff.monitor.pending.limit, default 10,000) and releases the lock every dfs.namenode.decommission.backoff.monitor.pending.blocks.per.lock blocks (default 1,000). It exists because the default monitor can flood the queues and stall the NameNode when many dense nodes are removed at once.

How long should it take

Work it out before you start. A DataNode holding 40 TB in 128 MB blocks carries about 300,000 blocks. Each block needs one new replica, so the cluster must copy 40 TB.

Two limits bound the rate. The scheduling limit is the redundancy monitor: each run schedules roughly dfs.namenode.replication.work.multiplier.per.iteration (default 2) blocks per live DataNode. With 100 live nodes that is about 200 blocks every three seconds, or roughly 67 blocks and 8.5 GB per second at most. The transfer limit is per node: a DataNode works on only a few replication streams at once (dfs.namenode.replication.max-streams, default 2, with a hard limit of 4), and the decommissioning node's own disks and network carry a large share of the reads. A single node pushing 300 to 500 MB/s is a realistic ceiling, which puts 40 TB at roughly 22 to 37 hours if it is the main source.

That arithmetic explains common tuning moves. Decommission several nodes in different racks together, since sources spread out and the total rate rises. Raise the work multiplier and max streams cautiously during quiet hours, watching client latency, because replication competes with jobs for disk and network. And never expect a dense node to drain in minutes.

SettingDefaultRaise it when
dfs.namenode.replication.work.multiplier.per.iteration2replication queue is long and network idle
dfs.namenode.replication.max-streams2per-node transfer is the bottleneck
dfs.namenode.decommission.blocks.per.interval500000rarely; larger scans hold the lock longer
dfs.namenode.decommission.max.concurrent.tracked.nodes100decommissioning very many nodes at once

Watching progress

Progress is visible per node in the NameNode JMX bean Hadoop:service=NameNode,name=NameNodeInfo, whose DecomNodes attribute is a JSON string mapping each decommissioning node to counters: under-replicated blocks, blocks for which this node holds the only replica, and under-replicated blocks in open files. The last two are the ones that predict a stuck node.

import json, requests

NN = "http://nn1.example.com:9870"
bean = requests.get(f"{NN}/jmx",
                    params={"qry": "Hadoop:service=NameNode,name=NameNodeInfo"}).json()["beans"][0]
for node, s in json.loads(bean["DecomNodes"]).items():
    print(f"{node:40} under={s['underReplicatedBlocks']:>8} "
          f"only_here={s['decommissionOnlyReplicas']:>6} open_files={s['underReplicateInOpenFiles']:>6}")

Plot underReplicatedBlocks over time. A steady downward slope means healthy progress; a flat line well above zero means something structural is blocking it, not slowness.

Watch the cost side at the same time. Cluster-wide, the NameNode's count of under-replicated blocks should rise when the decommission starts and fall steadily with it; a rise that keeps growing means other nodes are failing or the queue is being refilled faster than it drains. On the DataNodes, compare network and disk utilisation on the leaving node with client read latency for jobs running on the cluster. If latency climbs, lower the work multiplier rather than pausing the decommission, because a paused node still carries data you are counting on removing.

Why decommissions get stuck

  • Open files. A block under construction cannot be safely re-replicated until the writer closes it. Long-lived writers, such as streaming ingest or HBase write-ahead logs, keep blocks on the node indefinitely. Find them with hdfs dfsadmin -listOpenFiles -blockingDecommission and roll or close those files.
  • Not enough eligible targets. A file with replication factor 10 on a cluster that will have 9 nodes left can never be satisfied. Lower the replication factor of such files with hdfs dfs -setrep.
  • Placement constraints. With rack awareness, if the leaving node is the last one in a rack, blocks that need a second rack may have nowhere valid to go.
  • Erasure-coded stripes. An RS-6-3 block group needs nine distinct DataNodes, and placement prefers distinct racks. Shrinking a small cluster below that width makes reconstruction impossible; see HDFS erasure coding. EC blocks are also reconstructed by decoding, which costs CPU and reads from several nodes.
  • Corrupt replicas elsewhere. If the other replicas of a block are corrupt, the leaving node holds the only good one and the copy must come from it. Check hdfs fsck / -list-corruptfileblocks before starting.
  • Lost exclude state. After a NameNode restart or failover, the exclude or JSON file on that NameNode determines the state. If it was not updated, the node is quietly back in service.

Maintenance mode: the cheaper option for short outages

If a node is going down for a kernel patch or a disk swap and coming back within hours, decommissioning wastes a full copy of its data. Maintenance mode tells the NameNode the node will return. In ENTERING_MAINTENANCE the NameNode only ensures every block on it has at least dfs.namenode.maintenance.replication.min live replicas elsewhere (default 1). Once that holds, the node is IN_MAINTENANCE and can be stopped; while it is down its blocks are not re-replicated. When the expiry time passes, or you remove the state, it returns to normal.

The price is reduced durability during the window. With the default minimum of 1, a block replicated three times may have a single live copy while the node is down. Raise the minimum to 2 for data you cannot afford to lose to a second failure, and keep maintenance windows short. Use decommissioning for permanent removal, hardware returns and any outage with an unknown end.

Maintenance also interacts with rolling work. Putting a whole rack into maintenance at once is allowed, but with rack-aware placement many blocks keep two of their three replicas in one rack, so the minimum must be satisfied first by copying, and the window starts later than expected. Move one rack at a time, wait for every node to reach IN_MAINTENANCE before stopping anything, and set the expiry long enough to cover the work plus a margin, because an expired node that is still down is treated as dead and triggers full re-replication.

Runbook

  1. Check health first: no missing or corrupt blocks in hdfs fsck /, and no files with replication factors above the remaining node count.
  2. Check capacity: the remaining nodes must hold the leaving node's data at your target utilisation.
  3. If YARN NodeManagers share the host, start a graceful YARN decommission so running containers finish: add the host to the YARN exclude file and run yarn rmadmin -refreshNodes -g 3600 -client.
  4. Add the host to the exclude or JSON file on every NameNode and run hdfs dfsadmin -refreshNodes.
  5. Watch the per-node counters; list and resolve blocking open files.
  6. When the node reports Decommissioned, stop the DataNode, then remove it from the include and exclude files and refresh again.
  7. Run the balancer if the remaining nodes are now uneven.

What to do next

  1. Compute the expected duration for your densest node using the arithmetic above, and write it into the runbook.
  2. Add an alert when a decommissioning node's under-replicated count stays flat for more than an hour.
  3. Switch to the JSON hosts provider and manage it from configuration tooling for both NameNodes.
  4. Inventory long-lived open files, such as ingest logs and HBase WALs, and agree how to roll them during a decommission.
  5. Check that the cluster keeps enough nodes and racks for its widest erasure-coding policy after removal.
  6. Rehearse maintenance mode on one node with a short expiry, and decide your maintenance replication minimum.
Key takeaway: Decommissioning a DataNode means letting the NameNode re-replicate every block the node holds while it keeps serving reads, then stopping it only when it reports Decommissioned. The admin monitor finds short blocks and the redundancy monitor copies them, throttled by the work multiplier and per-node stream limits, so dense nodes take hours. Stuck decommissions almost always come from open files, impossible replication factors, placement or erasure-coding constraints, or exclude files out of sync between NameNodes. For short outages, maintenance mode avoids the full copy at a durability cost.