Every Hadoop service is a Java process: the NameNode, DataNodes, JournalNodes and ZKFC failover controllers in HDFS, the ResourceManager, NodeManagers and history servers in YARN. Their heaps, garbage collectors and pauses decide whether the cluster is stable. A NameNode that pauses for a minute can trigger a failover; a ResourceManager that runs out of heap takes every running application with it. Most production Hadoop outages that look like network or ZooKeeper problems turn out to be a JVM that stopped for too long.
This page is about those long-running service JVMs, not task containers. It explains what fills each daemon's heap, where Hadoop 3 reads the settings, how to read the pause monitor in the logs, and how a pause becomes an outage through specific timeouts. Collector flags by JDK are covered in JVM tuning for Hive and apply equally here; container and task heaps are in YARN containers.
The daemons and what fills each heap
Heap needs differ by orders of magnitude between daemons, so tune each one separately rather than with one global value:
| Daemon | What fills the heap | Grows with | Typical sensitivity |
|---|---|---|---|
| NameNode | Namespace: inodes, blocks, block locations | Files, directories and blocks | Very high: the whole namespace lives in memory |
| DataNode | Replica map, transfer threads, buffers | Replicas per node, concurrent transfers | Moderate |
| JournalNode | Edit log batches in flight | Edit rate | Low, but pauses stall NameNode writes |
| ZKFC | ZooKeeper session and health checks | Nothing much | Low heap, but pause-sensitive |
| ResourceManager | Applications, attempts, containers, scheduler state | Running and retained completed apps, nodes | High on large or busy clusters |
| NodeManager | Container tracking, aux services such as the shuffle handler | Containers per node, shuffle load | Moderate |
| MapReduce JobHistory server | Cached job histories | Jobs loaded into its cache | Moderate |
The ResourceManager keeps a bounded list of completed applications in memory; yarn.resourcemanager.max-completed-applications defaults to 1000 in the current yarn-default.xml. Raising it for longer history in the UI costs heap on a busy cluster, and a timeline or history service is the better place for long history.
Where Hadoop 3 reads heap and JVM options
Hadoop 3 rewrote the shell scripts and renamed the per-daemon variables. Settings live in hadoop-env.sh for HDFS and common options and in yarn-env.sh for YARN daemons:
# hadoop-env.sh
# Default heap for every Hadoop JVM; no unit means MB. There is no default:
# without it the JVM sizes the heap from machine memory.
export HADOOP_HEAPSIZE_MAX=4g
export HADOOP_HEAPSIZE_MIN=4g
# Per-daemon options are appended to HADOOP_OPTS; an -Xmx here wins.
export HDFS_NAMENODE_OPTS="-Xms64g -Xmx64g -XX:+UseG1GC -XX:MaxGCPauseMillis=200 \
-XX:+ParallelRefProcEnabled -XX:+AlwaysPreTouch \
-XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=/var/crash/hadoop \
-Xlog:gc*:file=${HADOOP_LOG_DIR}/gc-namenode.log:time,uptime,level,tags:filecount=10,filesize=100M"
export HDFS_DATANODE_OPTS="-Xms4g -Xmx4g -XX:+UseG1GC"
export HDFS_JOURNALNODE_OPTS="-Xms1g -Xmx1g"
export HDFS_ZKFC_OPTS="-Xms512m -Xmx512m"
# yarn-env.sh
export YARN_RESOURCEMANAGER_HEAPSIZE=16g # defaults to HADOOP_HEAPSIZE_MAX
export YARN_NODEMANAGER_HEAPSIZE=2g
export YARN_RESOURCEMANAGER_OPTS="-XX:+UseG1GC -XX:+HeapDumpOnOutOfMemoryError"Three details from the shipped templates and scripts matter. A heap value without a unit is read as megabytes. An -Xmx inside a daemon's _OPTS takes precedence over HADOOP_HEAPSIZE_MAX. And the Hadoop 2 names such as HADOOP_NAMENODE_OPTS still work but are deprecated: the scripts print a warning that the old variable has been replaced by HDFS_NAMENODE_OPTS and use the old value. After an upgrade, grep for that warning in daemon start-up output; a stale variable silently overriding the new one is a common source of confusion.
The unified -Xlog syntax above needs JDK 9 or later; Java 8 uses the older -Xloggc family, and mixing them stops the JVM from starting. The Hadoop Java Versions page on the Apache wiki states that Hadoop 3.3 and later support Java 8 and Java 11 at runtime only, and that from Hadoop 3.5 the server side requires JDK 17. Check the exact release you run before changing JDKs, and check every flag again after the change.
The NameNode heap
The NameNode heap is sized from the namespace: every file, directory and block is an object in memory, and planning guidance commonly cites about 1 GB of heap per million blocks as a starting point. HDFS performance tuning covers that estimate and the small-file problem behind it. Beyond the total, three rules are specific to the JVM.
Set -Xms equal to -Xmx and add -XX:+AlwaysPreTouch so the full heap is committed and touched at start-up. A NameNode that grows its heap under load pays for page faults and heap resizing at the worst moment, and a host that cannot provide the memory should fail at start, not hours later.
Size the standby identically. The standby NameNode holds the same namespace and must take over at full load; a smaller standby fails exactly when you need it. See NameNode high availability.
Prefer G1 for large heaps on current JDKs. Older guides recommend CMS, which was removed in JDK 14. With G1, large heaps holding hundreds of millions of small long-lived objects mostly stress mixed collections and remembered sets; start with the default pause target, enable GC logging, and change one thing at a time based on the logs. Very large arrays, such as those used while processing big block reports, can become humongous allocations; a larger -XX:G1HeapRegionSize is the lever if the GC log shows humongous allocations triggering collections.
Worker and YARN daemon heaps
The other daemons need far less heap than the NameNode, but each has its own driver, and the right size is the one that keeps heap used after collections comfortably below the maximum at peak load. Measure that from the GC log or the MemHeapUsedM metric over a busy week rather than copying a value from another cluster.
DataNodes hold a map of the replicas they store and a thread with buffers for every concurrent read or write. Heap grows with replicas per node and with the transfer thread count, so dense nodes with many disks need more than the defaults suggest. Short-circuit reads and some I/O paths use direct (off-heap) memory, which -Xmx does not cover; leave room for it when you budget host memory.
The ResourceManager grows with running applications, containers, retained completed applications and the number of nodes. Give the standby the same heap as the active, as with the NameNode, because it rebuilds the same state on failover.
NodeManagers are usually small, but the MapReduce or Spark shuffle service runs inside them as an auxiliary service, so heavy shuffle load raises their heap and direct-memory use. Remember that NodeManager heap is taken from the same host memory you advertise to YARN for containers: subtract every daemon heap on a worker before setting the memory YARN may hand out.
Reading the pause monitor
The NameNode, DataNodes, ResourceManager and NodeManagers each run a JvmPauseMonitor thread that sleeps for 500 ms in a loop and measures how much longer than that the sleep actually took. Extra time means the whole process was stopped. It logs at INFO when the extra time exceeds jvm.pause.info-threshold.ms (default 1000) and at WARN above jvm.pause.warn-threshold.ms (default 10000). The message looks like this:
WARN org.apache.hadoop.util.JvmPauseMonitor: Detected pause in JVM or host machine (eg GC): pause of approximately 14215ms
GC pool 'G1 Old Generation' had collection(s): count=1 time=13870ms
INFO org.apache.hadoop.util.JvmPauseMonitor: Detected pause in JVM or host machine (eg GC): pause of approximately 3412ms
No GCs detectedRead the second line first. If it names a GC pool whose time roughly matches the pause, the collector stopped the process: look at the GC log for why (heap too small, promotion failure, humongous allocations). If it says No GCs detected, the JVM did not cause the pause; the host did. Typical causes are swapping, a hypervisor stealing CPU, transparent huge page compaction, a stalled disk holding a log write, or cgroup CPU throttling. Tuning GC will not fix those; disable swap or set vm.swappiness low on master hosts, check CPU steal and throttling counters, and look at the kernel log at the same timestamp.
The same information is exported as metrics in the JvmMetrics record: GcCount, GcTimeMillis, GcNumInfoThresholdExceeded, GcNumWarnThresholdExceeded and GcTotalExtraSleepTime, along with heap usage such as MemHeapUsedM. Alert on any increase of the warn counter on master daemons and graph heap used after collections; Hadoop monitoring covers collecting them.
How a pause becomes an outage
A pause becomes an outage when it outlasts a timeout somewhere else. The defaults in the current configuration files set the thresholds:
| If this pauses | Longer than | Then | Setting |
|---|---|---|---|
| ZKFC process, or the whole master host | 10 s | ZooKeeper session expires, the lock moves, failover starts | ha.zookeeper.session-timeout.ms = 10000 |
| Active NameNode JVM only | 45 s | ZKFC health RPC times out, ZKFC quits the election, failover starts | ha.health-monitor.rpc-timeout.ms = 45000 |
| A majority of JournalNodes, or the network to them | 20 s | Edits cannot be persisted; the NameNode shuts itself down | dfs.qjournal.write-txns.timeout.ms = 20000 |
| DataNode | 10.5 min | Marked dead, its replicas re-replicated | 2 x dfs.namenode.heartbeat.recheck-interval + 10 x dfs.heartbeat.interval |
| NodeManager | 10 min | Marked lost, its containers rescheduled | yarn.nm.liveness-monitor.expiry-interval-ms = 600000 |
Note which process paused. A GC pause in the NameNode JVM alone does not stop the separate ZKFC process, so its ZooKeeper session survives; the failover comes only when the health RPC gives up after 45 seconds. A pause of the whole master host, such as swapping, stops both and can cost a failover after about 10 seconds. If the old active wakes up it is fenced and cannot write, so clients see a burst of errors while they find the new active. Raising timeouts trades faster failure detection for fewer false failovers; fix the pause first and treat timeout changes as a last resort, made identically on every node.
Worked example: a failover every few days
A cluster reports a NameNode failover every few days, always mid-morning. The NameNode log shows a WARN from the pause monitor of about 50 seconds, naming the G1 old-generation pool, and the ZKFC log shows its health check of the NameNode failing in the same minute. The configuration has -Xmx48g with no -Xms, and the namespace has grown to 52 million blocks since the heap was last set.
The GC log shows heap used after each mixed collection rising to about 44 GB of 48: the heap is nearly full of live data, so G1 runs back-to-back mixed collections and finally a full collection. The planning figure of 1 GB per million blocks suggests about 52 GB before headroom, so the heap was already undersized by the rule of thumb as well as by the logs. The fix is a heap of 80 GB with -Xms equal to -Xmx and AlwaysPreTouch, applied first to the standby, which then takes over through a planned failover so the other node can be restarted with the same settings. Afterwards, heap used after collections settles near 55% and the warn counter stops rising. A second action, enforcing quotas and compacting small files, slows the growth that caused the problem.
Failure modes
| Symptom | Likely cause | Action |
|---|---|---|
| Pause WARN naming an old-generation pool | Heap too small for live data | Raise heap; reduce objects (small files, completed apps) |
| Pause with No GCs detected | Swap, CPU steal, throttling, disk stall | Fix the host; GC tuning will not help |
| JVM fails to start after a JDK upgrade | Removed GC or logging flags | Replace with unified logging and supported flags |
| Deprecation warning at start-up | Hadoop 2 variable names | Move values to HDFS_ or YARN_ variables |
| Standby fails during takeover | Standby heap smaller than active | Keep heap settings identical on both |
| ResourceManager OutOfMemoryError | Too many retained apps or nodes | Lower retained apps; raise heap; check heap dump |
| DataNode marked dead with no crash | Pause longer than 10.5 minutes or stuck host | Check pause log and host health |
What to do next
- List every daemon on each host with its configured -Xms, -Xmx and collector, taken from the running process, not the config files.
- Move any Hadoop 2 variable names to their Hadoop 3 equivalents and remove the deprecation warnings.
- Set -Xms equal to -Xmx and AlwaysPreTouch on NameNodes, and keep active and standby identical.
- Enable GC logging with rotation and heap dumps on OutOfMemoryError for every master daemon.
- Alert on GcNumWarnThresholdExceeded and on heap used after collection above about 70% on masters.
- For each pause, read whether a GC pool is named; fix the host when it is not.
- Disable swap or set swappiness low on master hosts and check CPU throttling for daemons in cgroups.
- Re-check NameNode heap against block count every quarter, and before any JDK change read the Java versions page for your release.