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.
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.
| Layer | What changes | Why it matters |
|---|---|---|
| On-disk layout | NameNode metadata layout version, DataNode storage layout version | If either changes, the old software cannot read the new files. You can roll back (losing new data) but not downgrade |
| Wire protocols | RPC between clients, NameNodes, DataNodes, ResourceManager, NodeManagers | Mixed versions must talk to each other for the whole upgrade window |
| Client APIs | Java APIs, CLI flags, deprecated configuration keys | Jobs compiled or scripted against the old version can break even when the cluster is healthy |
| Configuration and defaults | Default ports, renamed keys, changed defaults | A daemon that starts cleanly can listen on a port your firewall, monitoring or clients do not expect |
| Ecosystem | Hive, HBase, Spark, Oozie, Ranger, connectors | Each has its own supported Hadoop range and sometimes its own schema upgrade |
| Runtime | JDK and OS packages | A 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.
| Pattern | How it works | Downtime | Rollback | Fits when |
|---|---|---|---|---|
| Express (in place) | Stop everything, install new software, start the NameNode with -upgrade, start DataNodes, finalize later | Hours, planned | Until finalize: hdfs namenode -rollback, losing writes made since the upgrade | Layout changes, major jumps, clusters that can take a window |
| Rolling | Prepare a rollback image, upgrade standby then active NameNode, then DataNodes in batches, finalize | None for HDFS reads and writes | Downgrade if layouts unchanged; otherwise rollback with downtime | Minor and patch releases on an HA cluster |
| Side-by-side | Build a new cluster on the new version, copy data and metadata, shadow-run jobs, cut over | Minutes at cutover | Point clients back at the old cluster | Major jumps, hardware refresh, or when rollback must be instant |
| Canary | Run the new version on a small cluster with real workloads before touching production | None | Not needed; it is a rehearsal | Always, as a prefix to one of the others |
| Client-staged | Upgrade servers first, then roll clients and job jars over weeks | None | Per-team, by pinning the old client | After any server upgrade, when many teams own clients |
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.
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.
- Servers before clients. Upgrade the daemons, let them soak, then move gateways, edge nodes and job jars.
- 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.
- HDFS before YARN, YARN before applications. MapReduce and Spark jobs depend on both.
- Metastore schema before Hive services. Hive's
schematoolupgrades the metastore database schema; run it with-dryRunfirst, after a database backup, and never let two Hive versions write to one schema. - 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 productionThree 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.
| Moment | What you can still do | What it costs |
|---|---|---|
Before -rollingUpgrade prepare or -upgrade | Abort | Nothing |
| Rolling, layouts unchanged, not finalized | Downgrade | Software only; data written since is kept |
| Rolling or express, layout changed, not finalized | Rollback | Downtime, and every write since the upgrade started is lost |
| After finalize | Nothing in place | Restore from backup or rebuild |
| Side-by-side, before old cluster is retired | Repoint clients to the old cluster | Writes 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.
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
- List the six layers for your target version: layout versions, protocols, client APIs, defaults and ports, ecosystem ranges, JDK.
- Write down each ecosystem component's supported Hadoop range and take the intersection as your target.
- 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.
- Pick the pattern with the two-question test, then add a canary in front and a client-staged phase behind.
- Put the preflight script in front of every window and run it on the canary first.
- Write the rollback points into the runbook with exact commands and a soak period before finalize.
- Read the rolling upgrade runbook if your path is rolling, or DistCp with snapshots if it is side-by-side.