Hadoop 3.0 shipped in late 2017 and was the first major version change in years. It is not one feature but a bundle: some changes save storage, some remove old limits, and some break scripts, firewall rules and dependency trees on the day you upgrade. The 3.x line has kept moving since then; 3.4.3 was released in February 2026.

This article maps what changed, why each change exists, what it costs to adopt and how to use it, with commands and a worked example. It stays at the level of a decision map; the deeper pages are linked as each feature comes up. If you are planning the upgrade itself, the procedure is in Hadoop upgrade patterns.

Advertisement

The change map

Hadoop 3.x changes by layerHDFSEC, 3+ NameNodes, disk balancerrouter federation, new portsYARNresource types, GPUsopportunistic, Timeline v2Clients and toolsshaded jars, new shellJava 8 baseline, heap autoStorage connectorsS3A on AWS SDK v2 (3.4), ABFS, ADLS, OSSRetired or revertedS3Guard removed; NN RPC back to 8020Your clusteradopt per feature: config, ports, clients, data layout
The Hadoop 3 line touched every layer. Some changes are switched on per directory or per queue, others change every client and firewall rule on day one.
ChangeProblem it solvesCost of adopting
Erasure coding3x storage overhead of replicationCPU and network on reads of lost data; bad for small files
More than two NameNodesOnly one failure tolerated in HAMore JournalNodes and config; small
Non-ephemeral default portsDaemons failing to bind on busy hostsEvery URL, firewall rule and monitor changes
Shaded client jarsGuava and Jackson conflicts in applicationsSwitch dependencies; mostly painless
Java 8 baselineOld language and runtimeUpgrade JVMs everywhere
Shell script rewriteInconsistent, buggy start scriptsCustom wrappers and env files need review
YARN resource typesOnly memory and vcores schedulableNode-level config and isolation work
Opportunistic containersIdle capacity between guaranteed containersJobs must tolerate being queued or killed
Intra-DataNode disk balancerSkew after replacing a diskRun it; small

Erasure coding

Replication stores three full copies of each block, so 1 PB of data takes 3 PB of disk. Erasure coding splits data into cells, stripes them across several data blocks and computes parity blocks, so any lost block can be rebuilt from the rest. With the Reed-Solomon policy RS-6-3-1024k, every 6 data blocks get 3 parity blocks: 9 blocks hold 6 blocks of data, a 1.5x overhead, and any 3 of the 9 can be lost. RS-10-4-1024k gives 1.4x with 4 parity blocks.

The costs are real. Reading a file whose block is missing means reading 6 other blocks across the network and decoding, so recovery is far more expensive than copying a replica. Data locality is lost because each file is spread over many nodes. A policy needs at least as many DataNodes as data plus parity blocks, and ideally as many racks, to keep rack-level fault tolerance. Files much smaller than one stripe waste space because parity is still written. EC has also lacked some write features that replication supports, such as hflush and hsync, so keep write-ahead and streaming data replicated.

hdfs ec -listPolicies                                   # see which policies exist and are enabled
hdfs ec -enablePolicy -policy RS-6-3-1024k
hdfs ec -setPolicy -path /warehouse/cold -policy RS-6-3-1024k
hdfs ec -getPolicy -path /warehouse/cold

# A policy applies to NEW files only. Rewrite existing data to convert it:
hadoop distcp -update /warehouse/cold_old/2024 /warehouse/cold/2024

The usual pattern is replication for hot and recently written data, erasure coding for cold, large, write-once files such as older partitions of a table. Policies, striping layout and the read path are covered in HDFS erasure coding.

Advertisement

More than two NameNodes

Hadoop 2 HA allowed exactly two NameNodes, active and standby. During maintenance on the standby you had no protection at all. Hadoop 3 lets you configure more than two; the documentation recommends three and advises against going beyond five because of the extra edit-log traffic. Pair three NameNodes with five JournalNodes and the cluster tolerates two failures of each.

<!-- hdfs-site.xml -->
<property><name>dfs.nameservices</name><value>prod</value></property>
<property><name>dfs.ha.namenodes.prod</name><value>nn1,nn2,nn3</value></property>
<property><name>dfs.namenode.rpc-address.prod.nn3</name><value>nn3.example.com:8020</value></property>
<property><name>dfs.namenode.http-address.prod.nn3</name><value>nn3.example.com:9870</value></property>
<property>
  <name>dfs.namenode.shared.edits.dir</name>
  <value>qjournal://jn1:8485;jn2:8485;jn3:8485;jn4:8485;jn5:8485/prod</value>
</property>

Clients find the active NameNode through the failover proxy provider, which tries each configured NameNode in turn, so every client configuration must list all three. The failover mechanics are explained in NameNode high availability. Router-based federation, added in the same era, puts a routing layer in front of several namespaces so clients see one mount table without client-side configuration; it solves namespace scale, not availability.

Default ports moved

Many Hadoop 2 defaults sat in the Linux ephemeral range (32768 to 61000), so a daemon could fail to start because an outgoing connection had already taken its port. Hadoop 3 moved them (HDFS-9427). KMS moved too (HADOOP-12811), because its old port clashed with the HBase Master.

Daemon and endpointHadoop 2Hadoop 3
NameNode web UI (HTTP)500709870
NameNode HTTPS504709871
Secondary NameNode HTTP / HTTPS50090 / 500919868 / 9869
DataNode data transfer500109866
DataNode IPC500209867
DataNode HTTP / HTTPS50075 / 504759864 / 9865
KMS160009600
NameNode RPC80208020 (briefly 9820 in 3.0.0, then reverted)

Grep every monitoring check, firewall rule, WebHDFS URL and runbook for the old numbers before the upgrade. The NameNode RPC port is the trap in the other direction: 3.0.0 moved it to 9820, and 3.0.1 moved it back to 8020 (HDFS-12990), so a configuration copied from a 3.0.0 example may carry a port nothing else uses.

Clients, Java and scripts

Applications that depend on hadoop-client inherit Hadoop's own versions of Guava, Jackson, Protobuf and others, which routinely clash with the application's versions. Hadoop 3 publishes shaded artifacts that relocate those dependencies inside the jar:

<dependency>
  <groupId>org.apache.hadoop</groupId>
  <artifactId>hadoop-client-api</artifactId>
  <version>3.4.3</version>
</dependency>
<dependency>
  <groupId>org.apache.hadoop</groupId>
  <artifactId>hadoop-client-runtime</artifactId>
  <version>3.4.3</version>
  <scope>runtime</scope>
</dependency>

Compile against the API jar only; the runtime jar carries the relocated third-party classes. Code that reaches into Hadoop's private classes or relies on Hadoop's transitive dependencies will fail to compile, which is the point.

Hadoop 3.0 raised the minimum Java version from 7 to 8. Later 3.3 releases added Java 11 as a supported runtime. Support for newer JDKs has been arriving through the 3.4 line; check the release notes of the exact version you deploy before running daemons on anything newer than 11.

The shell scripts were rewritten as well. Daemons are started with hdfs --daemon start namenode and yarn --daemon start resourcemanager instead of the old hadoop-daemon.sh family, HADOOP_HEAPSIZE is deprecated in favour of HADOOP_HEAPSIZE_MAX and HADOOP_HEAPSIZE_MIN, and the JVM can size the heap from host memory when neither is set. MapReduce task heaps can now be derived from the container size, so you no longer have to keep mapreduce.map.memory.mb and -Xmx in step by hand.

YARN: resource types, opportunistic containers, Timeline v2

Hadoop 2 YARN scheduled two resources: memory and vcores. Hadoop 3 generalised this into resource types, so a cluster can declare countable resources such as GPUs and schedule them like memory. Declare the type in resource-types.xml, have each NodeManager report its count, and let applications request it.

<!-- resource-types.xml on the ResourceManager and NodeManagers -->
<property><name>yarn.resource-types</name><value>yarn.io/gpu</value></property>

# Request one GPU for a distributed shell container
yarn jar hadoop-yarn-applications-distributedshell-*.jar \
  -jar hadoop-yarn-applications-distributedshell-*.jar \
  -shell_command nvidia-smi -container_resources memory-mb=4096,vcores=2,yarn.io/gpu=1

Scheduling a GPU is not the same as isolating it. Enable the GPU plugin and cgroups on the NodeManagers, or two containers that were each granted one GPU can still both see all of them.

Opportunistic containers are a second execution type. They are queued at the NodeManager and run when guaranteed containers leave resources idle, and they can be preempted when guaranteed work needs the space. They raise utilisation for short, retry-tolerant tasks and should not be used for anything that cannot lose progress.

Timeline Service v2 rebuilt the application history store on a scalable backend (HBase) with distributed collectors. It shipped as an alpha in 3.0 and matured in later releases; confirm its status in your version's documentation before relying on it for production history.

The intra-DataNode disk balancer

HDFS's cluster balancer moves blocks between DataNodes. It cannot fix skew inside one DataNode, which is exactly what happens when you replace a failed disk: the new disk is empty while its neighbours are 85 percent full, and the round-robin volume choice keeps writing to all of them equally. Hadoop 3 adds the intra-DataNode disk balancer, which plans and runs moves between volumes of one node.

hdfs diskbalancer -plan dn17.example.com           # writes a JSON plan, prints its path
hdfs diskbalancer -execute /system/diskbalancer/<date>/dn17.example.com.plan.json
hdfs diskbalancer -query dn17.example.com          # progress

It must be enabled with dfs.disk.balancer.enabled on the DataNodes; check the default in your version. Bandwidth controls and planning details are in the HDFS disk balancer.

What changed after 3.0

The 3.x line kept changing after 3.0, mostly in the cloud connectors. S3Guard, a DynamoDB-backed metadata layer that made S3 listings consistent for S3A, became unnecessary once Amazon S3 became strongly consistent; it was removed from the source (HADOOP-17409), and in 3.4.0 S3A instances configured with a real S3Guard metastore fail to start. The same release moved S3A from AWS SDK for Java v1 to v2 (HADOOP-18073), which changes credential provider class names and some configuration, so test custom credential chains before upgrading. Azure gained the ABFS connector for ADLS Gen2, and the 3.0 release had already added Azure Data Lake and Aliyun OSS connectors.

The practical consequence: "Hadoop 3" names a decade of releases, not one feature set. Pin the exact version in every conversation about behaviour.

Worked example: a cold zone on erasure coding

A 60-node Hadoop 2 cluster in 6 racks stores 2.4 PB of raw data, which is 800 TB of logical data under replication factor 3. Monitoring shows that 500 TB of it is partitions older than 90 days that are read a few times a month.

  1. Upgrade first, unchanged. Move to 3.x with the existing replication everywhere. Update ports in monitoring and firewalls, switch application builds to the shaded client and add the third NameNode plus two more JournalNodes.
  2. Choose the policy for the topology. RS-6-3-1024k writes 9 blocks per stripe. With only 6 racks, losing one rack can take out 2 blocks of a stripe, which RS-6-3 still survives because it tolerates 3. Check whether your version offers a policy that suits your rack count before choosing a wider one such as RS-10-4.
  3. Do the arithmetic. 500 TB under replication uses 1,500 TB of disk. Under RS-6-3 it uses 750 TB, freeing 750 TB, nearly a third of the cluster.
  4. Convert in batches. Set the policy on new cold directories and use distcp to rewrite one month of partitions at a time, outside peak hours, then swap the paths in the metastore and delete the old copies. The rewrite reads 500 TB and writes 750 TB, so throttle it with distcp's bandwidth and map count options.
  5. Measure the cost. Watch reconstruction work and read latency on the converted directories for a month before converting more.

Failure modes

FailureSymptomFix
Old ports in monitoringDashboards blank after the upgradeInventory and change every 500xx reference
EC on small filesDisk use rises after conversionKeep small or hot files replicated; compact first
EC with too few racksData unavailable after a rack lossMatch data plus parity to the rack layout
Clients list only two NameNodesFailover to nn3 breaks some jobsShip the same hdfs-site.xml to every client
S3Guard settings left in core-site.xmlS3A fails to initialise on 3.4Remove the metastore settings
GPUs scheduled but not isolatedJobs collide on one deviceEnable the GPU plugin with cgroups

Trade-offs

Hadoop 3's storage savings are the main financial reason to upgrade, and they come with slower recovery and lost locality, so they suit cold data only. The availability, port and shell changes are pure operational hygiene: some one-off work in exchange for fewer problems later. Resource types matter if you schedule accelerators on YARN; if GPU work runs on Kubernetes instead, they matter little. Weigh all of it against the alternative of moving storage to object stores, where many of these HDFS features do not apply.

What to do next

  1. Write down the exact Hadoop version you run and the one you target, and read the release notes between them.
  2. Grep configs, monitors and firewall rules for 500xx ports and plan the change.
  3. Find your largest cold, large-file directories and compute the savings of an erasure coding policy that fits your rack count.
  4. Switch one application build to hadoop-client-api and hadoop-client-runtime.
  5. Add a third NameNode and five JournalNodes on a staging cluster and drill a double failure.
  6. Adopt router federation only if namespace size is actually your limit.
  7. Keep the upgrade itself separate from feature adoption: upgrade, soak, then turn features on one at a time.
Key takeaway: Hadoop 3 is a bundle of independent changes. Erasure coding halves the disk cost of cold data at the price of slower recovery, extra NameNodes remove the single-failure limit, moved ports and the shell rewrite change every runbook, shaded jars end dependency clashes, and YARN can schedule GPUs. Upgrade first, then adopt each feature deliberately, with the exact 3.x version pinned.