A Hadoop cluster is rarely allowed to stop. Nightly pipelines, interactive queries and streaming jobs all read HDFS, so the classic upgrade, stop everything, install, start with -upgrade and finalize later, costs an outage that many teams cannot schedule. A rolling upgrade replaces the software one daemon at a time while the cluster keeps serving, supported for HDFS since Hadoop 2.4.0.

It works only because of three mechanisms: NameNode high availability so one NameNode can be restarted while the other serves, a special fsimage written before anything changes so the cluster can still go back, and a DataNode shutdown that tells writing clients to wait instead of failing. This article explains each mechanism, gives the command sequence, shows how YARN fits in, and ends with a runbook for a 200-node cluster.

Advertisement

Why rolling upgrades need special support

Restarting daemons one at a time sounds simple. The difficulty is that for a while the cluster runs two software versions at once, and the new version may write metadata the old one cannot read. HDFS has two kinds of on-disk metadata with independent version numbers: the NameNode's namespace, stored as an fsimage plus edit logs, and each DataNode's block storage directory. Each carries a layout version. If a new release changes a layout, the old software cannot read it any more.

A normal upgrade handles this by keeping a snapshot of the old directories and allowing rollback until finalize. A rolling upgrade cannot stop the NameNode to do that, so it takes a different approach. Before any daemon is touched, the active NameNode writes a rollback fsimage: a checkpoint of the namespace at the moment the upgrade started. DataNodes, for their part, stop deleting block files during the upgrade. When the NameNode tells them to delete a block, they move it to a per-block-pool trash directory instead, so the pre-upgrade data is still on disk if you roll back. Both are released only when you finalize.

The price is disk. Every block deleted during the upgrade window stays on disk, so a long soak on a cluster that churns temporary data can fill DataNodes. Watch capacity during the window, and do not leave a cluster unfinalized for weeks.

HDFS rolling upgrade: state machine and component order1. Preparerollback fsimage2. Standby NNupgrade, restart3. Failoverthen old active4. DataNodesbatch by batch5. Soakold and new mixed6. Finalizedrop rollback imageRollbackdowntime, pre-upgrade stateDowngradeDNs first, then NNshealthydata problemsoftware bug, same layoutRollback and downgrade are only possible before Finalize
The rolling-upgrade state machine. Rollback returns data and software to the moment of Prepare; downgrade keeps data written since then but works only if no layout changed.

Step 1: prepare

Start from a healthy cluster: both NameNodes up, no missing blocks, no DataNodes dead for unknown reasons. Upgrade the software packages in staging first and read the release notes for layout-version changes, because they decide whether downgrade is possible.

# 0. Pre-flight: healthy namespace, no ongoing upgrade, both NameNodes up
hdfs haadmin -getAllServiceState
hdfs dfsadmin -report | head -20          # live/dead DataNodes, remaining capacity
hdfs fsck / | tail -5                     # no missing or corrupt blocks

# 1. Prepare: the active NameNode writes an fsimage for rollback
hdfs dfsadmin -rollingUpgrade prepare
# poll until the output says "Proceed with rolling upgrade"
until hdfs dfsadmin -rollingUpgrade query | grep -q "Proceed with rolling upgrade"; do
  sleep 30
done

The query command reports progress while the rollback image is being written. Do not continue until it says to proceed. The time taken grows with namespace size, since it is a full checkpoint. On a namespace with hundreds of millions of files, expect minutes. The mechanics are the same as a normal checkpoint, described in the checkpoint and JournalNode article.

Advertisement

Step 2: the NameNodes, standby first

With HDFS high availability, the upgrade never takes the namespace offline. Upgrade the standby, make it active, then upgrade the former active:

# 2. NameNodes: nn1 is active, nn2 is standby
ssh nn2 'hdfs --daemon stop namenode'
# ... install new Hadoop version on nn2 ...
ssh nn2 'hdfs --daemon start namenode -rollingUpgrade started'
hdfs haadmin -getServiceState nn2         # wait for "standby" and edit-log catch-up

hdfs haadmin -failover nn1 nn2            # nn2 (new version) becomes active

ssh nn1 'hdfs --daemon stop namenode'
# ... install new Hadoop version on nn1 ...
ssh nn1 'hdfs --daemon start namenode -rollingUpgrade started'

The -rollingUpgrade started startup option tells the new NameNode that a rolling upgrade is in progress, so it keeps the rollback image and understands the edit log written by the old version. Before failing over, confirm the upgraded standby has caught up on edits and is healthy; a failover to a NameNode that is still loading its fsimage is an outage. If you use automatic failover through ZooKeeper, the failover command still works, but make sure the ZKFC on each host is running the version you expect. The NameNode HA article covers fencing and the failover controller.

Without HA the procedure still exists, but it is not free: stop the Secondary NameNode, then stop, upgrade and start the NameNode with -rollingUpgrade started, then upgrade the Secondary. The namespace is unavailable while the NameNode restarts, and clients retry or fail during that gap. In a federated cluster, prepare and finalize run once per namespace, and each active and standby pair is upgraded as above.

Step 3: the DataNodes, batch by batch

DataNodes are upgraded in small groups. The key command is hdfs dfsadmin -shutdownDatanode <host:ipc_port> upgrade. The upgrade argument tells the DataNode it is restarting soon, so before it stops it sends clients in an active write pipeline an out-of-band signal. Those clients pause and wait for the node to return instead of immediately rebuilding the pipeline without it. How long they wait is set by dfs.client.datanode-restart.timeout, 30 seconds by default. A restart that takes longer is treated as a normal DataNode failure, and the pipeline recovers without the node.

#!/usr/bin/env bash
# 3. DataNodes, one rack-aware batch at a time. IPC port default 9867 (dfs.datanode.ipc.address).
set -euo pipefail
BATCH_FILE=$1                             # hostnames, all from ONE rack

for dn in $(cat "$BATCH_FILE"); do
  hdfs dfsadmin -shutdownDatanode "$dn:9867" upgrade
done

for dn in $(cat "$BATCH_FILE"); do        # wait until each DataNode is really down
  while hdfs dfsadmin -getDatanodeInfo "$dn:9867" >/dev/null 2>&1; do sleep 2; done
  ssh "$dn" 'install-new-hadoop && hdfs --daemon start datanode'
done

for dn in $(cat "$BATCH_FILE"); do        # and back before the next batch
  until hdfs dfsadmin -getDatanodeInfo "$dn:9867" >/dev/null 2>&1; do sleep 5; done
done
hdfs dfsadmin -report | grep -E "Live datanodes|Dead datanodes"

Batch selection matters more than batch speed. With replication factor 3 and rack-aware placement, every block has replicas on at least two racks, so taking down several nodes from one rack at a time never makes a block unavailable. Taking down nodes from several racks at once can. The rack awareness article explains the placement rule this relies on. If you use erasure coding, check the policy: a 6+3 Reed-Solomon stripe tolerates three missing cells, and stripes are spread across racks differently from replicas.

Also avoid triggering re-replication. The NameNode considers a DataNode dead only after roughly ten minutes without heartbeats with default settings, so a restart that takes a minute or two creates no replication storm. A batch that hangs for longer does. Keep an eye on under-replicated block counts, and stop the rollout if they climb.

Step 4: soak, then finalize, downgrade or roll back

After the last DataNode, the cluster runs entirely on the new version, but it is still in rolling-upgrade mode: the rollback image is kept and DataNodes still keep deleted blocks. This is the time to run real workloads and compare metrics. You then have three exits.

ExitWhenData written since PrepareDowntimeHow
FinalizeNew version is healthyKeptNonehdfs dfsadmin -rollingUpgrade finalize
DowngradeSoftware bug, no layout changeKeptNone (rolling)Reinstall old software: DataNodes first with shutdownDatanode ... upgrade, then NameNodes started normally; finalize
RollbackData corruption or layout change blocks downgradeLostFull clusterStop all and reinstall the old release; start NN1 as active with -rollingUpgrade rollback, run -bootstrapStandby on NN2 and start it normally, then start DataNodes with -rollback

The difference between downgrade and rollback is the most misunderstood part. Downgrade changes software only: data written since Prepare stays, and the cluster keeps serving. It works only if neither the NameNode nor the DataNode layout version changed, and DataNodes must be downgraded before NameNodes. Rollback restores the namespace to the rollback image and DataNodes to their pre-upgrade state. Every file created, appended or deleted since Prepare is undone, and it requires stopping the whole cluster. Treat rollback as disaster recovery, and tell data owners what it would cost before you start.

Finalizing deletes the rollback image and DataNode trash. After that, neither exit is available; the only way back is restoring from backup or a replica cluster.

YARN and the rest of the stack

The HDFS procedure does not cover YARN, JournalNodes or ZooKeeper. The Hadoop documentation notes that upgrading JournalNodes and ZooKeeper may incur cluster downtime, so plan them separately: JournalNodes can usually be restarted one at a time while a quorum remains, but verify that the new version reads the old edit-log segments in staging first.

YARN needs two features for restarts not to kill running containers. ResourceManager recovery stores application state, usually in ZooKeeper, so a restarted or failed-over ResourceManager resumes scheduling existing applications. NodeManager recovery stores container state on local disk so a restarted NodeManager reattaches to containers that kept running. The NodeManager's RPC port must be fixed, or it comes back on a different port and the ResourceManager cannot reconcile it.

<!-- yarn-site.xml: prerequisites for restarting YARN daemons without killing containers -->
<property><name>yarn.resourcemanager.recovery.enabled</name><value>true</value></property>
<property><name>yarn.resourcemanager.store.class</name>
  <value>org.apache.hadoop.yarn.server.resourcemanager.recovery.ZKRMStateStore</value></property>
<property><name>yarn.nodemanager.recovery.enabled</name><value>true</value></property>
<property><name>yarn.nodemanager.recovery.dir</name><value>/var/lib/hadoop-yarn/nm-recovery</value></property>
<property><name>yarn.nodemanager.recovery.supervised</name><value>true</value></property>
<!-- must be a fixed port: port 0 would change across the restart -->
<property><name>yarn.nodemanager.address</name><value>0.0.0.0:45454</value></property>

With those enabled, upgrade the standby ResourceManager, fail over, upgrade the other, then restart NodeManagers in rack batches. Client libraries and application jars are the last concern: jobs submitted during the window may run with old client jars against new servers. Hadoop keeps wire compatibility within a major line, but test the frameworks you run, such as Hive, Spark and HBase, against the new servers in staging. Major-version jumps deserve their own compatibility analysis rather than an assumption that the rolling procedure covers them.

Worked plan: 200 DataNodes across 10 racks

A cluster has two NameNodes, three JournalNodes, 200 DataNode and NodeManager hosts in 10 racks of 20, replication factor 3, and a 1.2 PB namespace with 180 million files. Planning figures:

  • Prepare: rollback image for 180 million files, measured in staging at about 6 minutes.
  • NameNodes: about 20 minutes each including fsimage load and edit catch-up; one failover.
  • DataNodes: batches of 5 hosts from one rack, 4 batches per rack. Each batch takes about 6 minutes: stop, install, start and re-register with block reports. 40 batches take about 4 hours, with no block ever losing more than one replica.
  • Soak: 48 hours of normal workload, watching DataNode capacity because deleted blocks accumulate in trash; 3% of capacity was the churn estimate for two days.
  • Finalize, confirm DataNode capacity drops back as the upgrade trash is released, and record the new version in the inventory.

The plan names a stop condition for each phase: under-replicated blocks rising for more than 15 minutes, a DataNode not returning within 10 minutes, or job failure rate above baseline. Any of them pauses the rollout; a data problem triggers the rollback decision while it is still possible.

Failure modes

  • Continuing before the rollback image is ready. Upgrading the standby NameNode before query says to proceed leaves no safe rollback point.
  • Multi-rack batches. Restarting nodes from several racks at once makes some blocks temporarily unreadable and fails reads.
  • Trash fills disks. A long soak on a cluster with heavy delete churn, such as temporary Hive tables, fills DataNode volumes. Shorten the soak or free space before starting.
  • Slow DataNode restarts. Restarts beyond the client restart timeout cause pipeline recovery; beyond the dead-node interval, re-replication. Pre-stage packages so the restart is only a process restart.
  • Downgrade across a layout change. Blocked by design; the only way back is rollback with downtime and data loss since Prepare. Read the layout-version notes before you start.
  • Ephemeral NodeManager port. Containers are orphaned on restart. Set a fixed port before the upgrade, which itself needs a restart.

What to do next

  1. Confirm prerequisites now: HDFS HA with a healthy standby, ResourceManager HA and recovery, NodeManager recovery with a fixed RPC port.
  2. Rehearse the full sequence in staging with production-like namespace size and record timing for Prepare, each NameNode and one DataNode batch.
  3. Read the release notes for NameNode and DataNode layout-version changes, and write down whether downgrade is possible.
  4. Script DataNode batches by rack, with the shutdownDatanode and getDatanodeInfo checks and a stop condition on under-replicated blocks.
  5. Agree the soak length and the rollback decision owner with data owners before Prepare, including what data would be lost.
  6. Plan JournalNode and ZooKeeper upgrades as separate changes, and finalize promptly once the soak is clean.
Key takeaway: A Hadoop rolling upgrade keeps HDFS and YARN serving while the software changes underneath. Prepare writes a rollback fsimage, the standby NameNode is upgraded and failed over to, DataNodes restart in single-rack batches while clients wait for them, and nothing is permanent until finalize. Downgrade keeps new data but only works without layout changes; rollback discards everything since Prepare and needs downtime. Enable YARN recovery beforehand, rehearse with real timings, watch trash and under-replication, and finalize once the soak is clean.