Most Hadoop upgrade disasters are not caused by a bad command. They are caused by choosing the wrong pattern: a rolling upgrade attempted across a layout change nobody checked, an in-place major upgrade on a cluster whose business could not take the downtime, or a side-by-side migration that ran for nine months because nobody planned the cutover. The commands are the easy part and are well documented. This article is about the decision that comes before them.

We take an upgrade apart into the layers it touches, compare five patterns, set out ordering rules, rehearse against real NameNode metadata and work through a 2.x to 3.x decision. The rolling procedure has its own page, Hadoop Rolling Upgrade, in depth; here it is one pattern among several.

Advertisement

What an upgrade actually changes

A Hadoop version bump is really six changes travelling together, and each one has its own way of going wrong. Listing them is the first step of any plan, because the riskiest layer decides the pattern.

LayerWhat changesWhy it matters
On-disk layoutNameNode metadata layout version, DataNode storage layout versionIf either changes, the old software cannot read the new files. You can roll back (losing new data) but not downgrade
Wire protocolsRPC between clients, NameNodes, DataNodes, ResourceManager, NodeManagersMixed versions must talk to each other for the whole upgrade window
Client APIsJava APIs, CLI flags, deprecated configuration keysJobs compiled or scripted against the old version can break even when the cluster is healthy
Configuration and defaultsDefault ports, renamed keys, changed defaultsA daemon that starts cleanly can listen on a port your firewall, monitoring or clients do not expect
EcosystemHive, HBase, Spark, Oozie, Ranger, connectorsEach has its own supported Hadoop range and sometimes its own schema upgrade
RuntimeJDK and OS packagesA new JDK can change TLS defaults, GC behaviour and Kerberos encryption types

Two of these are checkable facts rather than judgement calls. The release notes and the hdfs namenode -upgrade behaviour tell you whether the metadata layout changes, and the Apache rolling upgrade documentation states the rule plainly: a newer release can be downgraded to the pre-upgrade release only if neither the NameNode nor the DataNode layout version changed. The same page notes that rolling upgrade itself is supported only from Hadoop 2.4.0 onwards. Everything else in the table you discover by rehearsal, which is why rehearsal gets its own section below.

The default ports are the classic Hadoop 3 surprise. Many web and data-transfer ports moved out of the Linux ephemeral range: for example the NameNode web UI moved from 50070 to 9870, the DataNode data-transfer port from 50010 to 9866 and the DataNode web UI from 50075 to 9864. Anything that hard-codes the old numbers, such as monitoring probes, firewall rules or WebHDFS URLs in scripts, breaks while the cluster itself reports healthy. Pin the ports explicitly in configuration during the upgrade and move them later as a separate change.

The five patterns

Every upgrade we have seen is one of these patterns or a combination of them. They differ in downtime, hardware cost and, above all, in what a rollback costs you.

PatternHow it worksDowntimeRollbackFits when
Express (in place)Stop everything, install new software, start the NameNode with -upgrade, start DataNodes, finalize laterHours, plannedUntil finalize: hdfs namenode -rollback, losing writes made since the upgradeLayout changes, major jumps, clusters that can take a window
RollingPrepare a rollback image, upgrade standby then active NameNode, then DataNodes in batches, finalizeNone for HDFS reads and writesDowngrade if layouts unchanged; otherwise rollback with downtimeMinor and patch releases on an HA cluster
Side-by-sideBuild a new cluster on the new version, copy data and metadata, shadow-run jobs, cut overMinutes at cutoverPoint clients back at the old clusterMajor jumps, hardware refresh, or when rollback must be instant
CanaryRun the new version on a small cluster with real workloads before touching productionNoneNot needed; it is a rehearsalAlways, as a prefix to one of the others
Client-stagedUpgrade servers first, then roll clients and job jars over weeksNonePer-team, by pinning the old clientAfter any server upgrade, when many teams own clients
Choosing an upgrade pattern: two questions decide most of itDoes the HDFS layout change?NameNode or DataNode layout versionNo: rolling or downgradablesame major line, certified pathYes, or major version jumprollback needs downtimenoyesRolling upgradeprepare, NNs, DNs in batches, finalizeCan you afford a second cluster?hardware, licences, dual-run weeksExpress (in place)planned downtimeSide-by-sidecopy data, cut overnoyesCanary first, alwaysa small cluster on the new buildWhatever you pick, clients and ecosystem jars move last, after the servers are stable.
Most of the decision follows from two facts: whether the layout changes and whether you can afford a second cluster.

The patterns compose: a real plan usually reads canary, then rolling or express for the servers, then client-staged for everything else. The trade-off is simple to state. Express concentrates risk in one window you can rehearse precisely; rolling avoids downtime but stretches the mixed-version period; side-by-side doubles hardware for weeks but gives the cleanest rollback. Pick by the cost of downtime and of a failed rollback.

Advertisement

Ordering rules that hold for every pattern

Whatever the pattern, the order of components is constrained by who talks to whom. Newer servers are built to accept older clients within a supported range; newer clients talking to older servers are not guaranteed to work. That single asymmetry produces the rules.

  1. Servers before clients. Upgrade the daemons, let them soak, then move gateways, edge nodes and job jars.
  2. NameNodes before DataNodes. The NameNode must understand the DataNodes' reports; in a rolling upgrade the standby goes first so a failover lands on new software.
  3. HDFS before YARN, YARN before applications. MapReduce and Spark jobs depend on both.
  4. Metastore schema before Hive services. Hive's schematool upgrades the metastore database schema; run it with -dryRun first, after a database backup, and never let two Hive versions write to one schema.
  5. One change at a time. Do not combine a JDK change, an OS upgrade and a Hadoop upgrade in one window. When something breaks you need to know which one did it.

Before choosing a target Hadoop version, write down the supported Hadoop range for every component you run, from its release notes or your vendor's matrix. The target is the intersection of those ranges.

Rehearse against a copy of the real metadata

The highest-value hour in any upgrade is spent starting the new NameNode against a copy of your production metadata on a lab machine. The fsimage holds the whole namespace, so the upgrade's metadata conversion, its memory use and its startup time can all be measured without any data blocks. Fetch the latest checkpoint, copy the NameNode's name directory, and start the new version with the upgrade flag against the copy.

# On production: download the most recent fsimage (does not stop anything).
hdfs dfsadmin -fetchImage /backup/nn-$(date +%Y%m%d)

# On a lab host with the NEW version, as the hdfs user, with a STANDALONE
# non-HA config: no dfs.namenode.shared.edits.dir, no HA nameservice keys and
# no route to production JournalNodes, or the lab NameNode can fence the real one.
# dfs.namenode.name.dir points at a COPY of a recent name directory.
nohup hdfs namenode -upgrade > nn-upgrade.log 2>&1 &   # the daemon never returns
until hdfs dfsadmin -safemode get 2>/dev/null; do sleep 10; done
# Load time: from process start to the image-loaded lines in nn-upgrade.log.
hdfs dfs -count -q /               # namespace totals should match production

Three numbers come out of this: load and conversion time (your downtime estimate), heap afterwards, and whether namespace totals match. The second rehearsal is workload replay on the canary: run real jobs, including the oldest, and compare output row counts and checksums, not just exit codes.

A preflight gate you can run before every window

Upgrades should start from a known-good cluster. The script below refuses to continue if HDFS is in safe mode or has missing or corrupt blocks, and prints the rolling upgrade state for inspection. It only greps the summary lines these commands print and stops on anything unexpected.

#!/usr/bin/env bash
# preflight.sh: refuse to start an upgrade window from an unhealthy cluster.
set -euo pipefail
fail() { echo "PREFLIGHT FAIL: $*" >&2; exit 1; }

hdfs dfsadmin -safemode get | grep -q "Safe mode is OFF" || fail "NameNode in safe mode"

report=$(hdfs dfsadmin -report)
# 3.x indents these under "Replicated Blocks:", 2.x does not.
echo "$report" | grep -E "^[[:space:]]*(Missing blocks|Under replicated blocks|Blocks with corrupt replicas):" || true
echo "$report" | grep -qE "^[[:space:]]*Missing blocks: 0$" || fail "missing blocks"

hdfs fsck / -list-corruptfileblocks | grep -q "has 0 CORRUPT files" \
  || fail "corrupt files present"

# Inspect by eye: it must show no rolling upgrade in progress.
hdfs dfsadmin -rollingUpgrade query

# Liveness ping for each DataNode's IPC port (9867 by default on 3.x).
while read -r dn; do
  hdfs dfsadmin -getDatanodeInfo "$dn" >/dev/null || fail "DataNode $dn not answering"
done < datanodes.txt

hdfs dfsadmin -fetchImage "/backup/pre-upgrade-$(date +%s)"
echo "PREFLIGHT OK"

Summary wording can shift between releases, so run it on the canary first; a grep that never matches fails closed.

Rollback points: know which door is still open

Every pattern has moments after which a rollback becomes more expensive or impossible. Write them into the runbook with the exact command, and tell stakeholders which door you are standing at.

MomentWhat you can still doWhat it costs
Before -rollingUpgrade prepare or -upgradeAbortNothing
Rolling, layouts unchanged, not finalizedDowngradeSoftware only; data written since is kept
Rolling or express, layout changed, not finalizedRollbackDowntime, and every write since the upgrade started is lost
After finalizeNothing in placeRestore from backup or rebuild
Side-by-side, before old cluster is retiredRepoint clients to the old clusterWrites made only on the new cluster must be copied back or replayed

The common mistake is finalizing early because DataNodes keep deleted blocks for rollback until finalize and disk fills. Set a soak period tied to a business cycle, such as one month-end close, and plan free space for it.

Worked example: a 2.x cluster going to 3.x

Take a 300-node cluster on Hadoop 2.10 running HDFS, YARN, Hive and Spark, with nightly batch SLAs and a business that can accept one six-hour window per quarter. The target is a current 3.x release on the same hardware or new hardware.

First, the rolling question. Rolling upgrade between 2.x and 3.x was tracked upstream in HDFS-11096, which the Apache JIRA still shows as unresolved (a sub-task, HDFS-11188, restored 2.x as the minimum supported peer version in 3.0.0-alpha2). Some vendors certified their own 2 to 3 rolling paths; unless yours did, in writing, treat this jump as express or side-by-side. Second, the layout changes across the major version, so a downgrade is off the table and rollback means downtime and lost writes.

Third, the rehearsal. Loading the fetched image on a lab NameNode with the 3.x build gives a conversion and startup time; say it measures 40 minutes for the namespace plus the time for 300 DataNodes to report. That fits inside six hours, so express is feasible on the same hardware. If the hardware is due for refresh anyway, side-by-side becomes attractive: build the new cluster, copy with snapshot-based DistCp as described in Hadoop DistCp, restore and upgrade a copy of the metastore, shadow-run the nightly jobs for two weeks, and cut over in a window measured in minutes.

Side-by-side upgrade: two clusters, one cutoverOld cluster (2.x)still serving all jobsNew cluster (3.x)fresh install, new layout1. bulk DistCp from snapshot s12. incremental: snapshot diff s1 to s2Metastore copydump, restore, schema upgrade3. metastore schema upgraded on new sideShadow job runscompare outputs, not exit codesCutover windowfreeze writes, final diff, flip DNSOld cluster read-onlyrollback = flip DNS back
Side-by-side keeps the old cluster intact until the new one has proven itself; rollback is a DNS change.

Either way the plan ends client-staged: edge nodes and job jars move to 3.x clients team by team after the servers soak.

Failure modes

  • NameNode heap exhaustion on first start of the new version. Prevented by measuring heap during the metadata rehearsal and sizing before the window.
  • Rolling upgrade started across a layout change, so the expected downgrade path does not exist. Prevented by checking layout versions in the release notes and on the canary before choosing the pattern.
  • HA failover during the window lands on a NameNode with the wrong software. Follow the documented order, standby first, and confirm NameNode HA health before each step.
  • DataNodes run out of disk because finalize is delayed for weeks. Plan the soak length against free space.
  • Monitoring goes blind because probes hit the old default ports. Pin ports in configuration and test probes on the canary.
  • Hive metastore upgraded while an old HiveServer2 still writes to it. Stop every writer before schematool -upgradeSchema.
  • Side-by-side drift: the new cluster misses writes made during the dual-run window. Use snapshot diffs and a final freeze rather than one-off copies, and keep the old cluster read-only after cutover.
  • Edit log problems surface only at restart. Check checkpoint and JournalNode health and force a fresh checkpoint first.

What to do next

  1. List the six layers for your target version: layout versions, protocols, client APIs, defaults and ports, ecosystem ranges, JDK.
  2. Write down each ecosystem component's supported Hadoop range and take the intersection as your target.
  3. Fetch the current fsimage and start the new NameNode against a copy of it on a lab host; record load time, heap and namespace totals.
  4. Pick the pattern with the two-question test, then add a canary in front and a client-staged phase behind.
  5. Put the preflight script in front of every window and run it on the canary first.
  6. Write the rollback points into the runbook with exact commands and a soak period before finalize.
  7. Read the rolling upgrade runbook if your path is rolling, or DistCp with snapshots if it is side-by-side.
Key takeaway: An upgrade is six changes at once: on-disk layout, wire protocols, client APIs, defaults, ecosystem and runtime. Whether the layout changes and whether you can afford a second cluster decide between rolling, express and side-by-side, and a canary goes in front of all of them. Upgrade servers before clients and NameNodes before DataNodes, rehearse on a copy of the real fsimage, know which rollback door is still open at every step, and do not finalize until the soak period has proven the new version.