Every piece of work that runs on a Hadoop YARN cluster runs inside a container: every MapReduce task, every Spark executor, every Flink TaskManager and every ApplicationMaster. The word suggests Docker, but a YARN container is something narrower: a grant of memory and CPU on one node, plus the process tree a NodeManager starts, watches and cleans up against that grant.

The YARN overview covers how requests become containers and how to size nodes. This article stays on the node. It follows one container from allocation to cleanup: what the ApplicationMaster sends to launch it, how files reach the node, which executor starts the process and as which user, how the container is killed, what each exit status means, and how to debug one that failed.

Advertisement

Allocation versus process

A container starts life as a record in the ResourceManager. When a scheduler places a request, the RM creates a Container with an id, the node it lives on, a Resource (memory in MB, vcores, and any custom resource types such as GPUs), a priority and a container token. The token is signed by the RM and lets the NodeManager verify that this ApplicationMaster really was granted this much capacity on this node. Nothing runs yet.

The AM receives the allocation on its next allocate call and must launch it by contacting the NodeManager directly. If it does not start the container within yarn.resourcemanager.rm.container-allocation.expiry-interval-ms (600,000 ms, ten minutes, by default), the RM takes the allocation back. That is why a stalled AM can hold capacity idle for up to ten minutes before the scheduler reuses it.

Container ids encode their history: container_e03_1727700000000_0042_01_000007 is container 7 of attempt 1 of application 42, from an RM that started at the timestamp shown, with e03 being the RM epoch that increases on restarts with recovery. Container 1 of each attempt is normally the ApplicationMaster itself.

A YARN container: allocated by the RM, launched and policed by one NodeManagerResource ManagerallocatesApp MasterAMRMClient / NMClient1. Allocationid, node, size, token2. startContainerlaunch context3. LocalizeHDFS to local disk4. Launchexecutor + script5. Monitormemory, cgroupsProcess treejava, python ...6. Exit + cleanupstatus, logs, dirsheartbeatexit status on next allocateSteps 3 to 6 all happen on one NodeManager; the RM only learns the outcome through the NM heartbeat.
The container lifecycle. Allocation happens in the ResourceManager; localization, launch, monitoring and cleanup happen on one NodeManager.

The launch context

To start a container, the AM sends the NodeManager the allocated Container and a ContainerLaunchContext. The launch context holds everything the NodeManager needs: the commands to run, environment variables, local resources to download first, security tokens (for example HDFS delegation tokens, so the process can read data as the user), auxiliary service data such as the shuffle service's credentials, and ACLs for viewing and modifying the application.

The command is a shell command line. Two conventions matter. First, redirect stdout and stderr into ApplicationConstants.LOG_DIR_EXPANSION_VAR, a placeholder the NodeManager replaces with the container's log directory; otherwise the output is lost. Second, size the JVM heap below the container: the NodeManager polices the whole process tree, so a 4096 MB container running -Xmx4g will be killed once metaspace, thread stacks and direct buffers push it over. Leave 15 to 25 percent for non-heap memory.

// Inside an ApplicationMaster: ask for containers, then launch each one on its node.
AMRMClient<ContainerRequest> rm = AMRMClient.createAMRMClient();
rm.init(conf); rm.start();
rm.registerApplicationMaster(host, 0, "");

NMClient nm = NMClient.createNMClient();
nm.init(conf); nm.start();

Resource size = Resource.newInstance(4096, 2);            // MB, vcores
for (int i = 0; i < 10; i++) {
    rm.addContainerRequest(new ContainerRequest(size, null, null, Priority.newInstance(1)));
}

int launched = 0;
while (launched < 10) {
    AllocateResponse resp = rm.allocate(launched / 10.0f);  // also the AM heartbeat
    for (Container ctr : resp.getAllocatedContainers()) {
        ContainerLaunchContext ctx = ContainerLaunchContext.newInstance(
                localResources(),                            // jars, configs, archives
                Map.of("JAVA_HOME", "/usr/lib/jvm/java-17"),
                List.of("$JAVA_HOME/bin/java -Xmx3276m com.example.Worker"
                        + " 1>" + ApplicationConstants.LOG_DIR_EXPANSION_VAR + "/stdout"
                        + " 2>" + ApplicationConstants.LOG_DIR_EXPANSION_VAR + "/stderr"),
                null, tokens(), null);
        nm.startContainer(ctr, ctx);
        launched++;
    }
    for (ContainerStatus st : resp.getCompletedContainersStatuses()) {
        log.info("{} exited {} : {}", st.getContainerId(), st.getExitStatus(), st.getDiagnostics());
    }
    Thread.sleep(1000);
}
Advertisement

Localization: getting files onto the node

Before the process starts, the NodeManager downloads the container's local resources from a shared filesystem, usually HDFS, into local disk. Each LocalResource has a URL, a type, a visibility, a size and a modification timestamp. The NodeManager checks the size and timestamp against the source, and localization fails if the file has changed since the AM described it. Overwriting a jar in place on HDFS while jobs are submitting is a common cause of mysterious launch failures.

// A LocalResource names a file on a shared filesystem plus the size and modification
// time the NodeManager must find when it downloads it.
Path jar = new Path("hdfs:///apps/worker/worker-1.4.2.jar");
FileStatus st = fs.getFileStatus(jar);

LocalResource res = LocalResource.newInstance(
        URL.fromPath(jar),
        LocalResourceType.FILE,                 // FILE, ARCHIVE (unpacked) or PATTERN
        LocalResourceVisibility.PUBLIC,         // PUBLIC, PRIVATE or APPLICATION
        st.getLen(),
        st.getModificationTime());              // a mismatch fails localization

Map<String, LocalResource> localResources = Map.of("worker.jar", res);  // link name in cwd

Visibility decides who can share the cached copy. PUBLIC resources are downloaded once per node and shared by every user; the source file must be world-readable, including every parent directory. PRIVATE resources are shared by all applications of one user. APPLICATION resources are shared only by containers of one application and deleted when it finishes. Types are FILE, ARCHIVE (unpacked after download) and PATTERN.

Public resources are downloaded by the NodeManager itself; private and application resources are downloaded by a localizer running as the user. By default the node keeps the public and private caches at a target of yarn.nodemanager.localizer.cache.target-size-mb (10240 MB), cleaned every localizer.cache.cleanup.interval-ms (ten minutes), and the public localizer uses yarn.nodemanager.localizer.fetch.thread-count (4) download threads. When a large framework jar is marked APPLICATION, every application downloads it again on every node. Marking it PUBLIC at a versioned path turns thousands of downloads into one per node.

Directory layout and launch_container.sh

Each local directory in yarn.nodemanager.local-dirs gets the same structure, and the NodeManager spreads containers across directories and disks. The working directory contains symlinks named after the keys of the local-resources map, so the process opens worker.jar without knowing where the cache is.

${yarn.nodemanager.local-dirs}/
  filecache/                                   PUBLIC resources, shared by all users
  usercache/alice/
    filecache/                                 PRIVATE resources, shared by alice's apps
    appcache/application_1727700000000_0042/
      filecache/                               APPLICATION resources for this app
      container_e03_1727700000000_0042_01_000007/   working directory (cwd)
        launch_container.sh                    generated: env, symlinks, your command
        worker.jar -> ../../filecache/...      symlink named by the map key
        tmp/
${yarn.nodemanager.log-dirs}/application_1727700000000_0042/
  container_e03_1727700000000_0042_01_000007/  stdout, stderr, application logs

The NodeManager writes launch_container.sh into the working directory: it exports the environment, including variables YARN adds such as CONTAINER_ID, NM_HOST, LOCAL_DIRS and LOG_DIRS, creates the symlinks and then runs your command. When a container fails before your code logs anything, this script is the first thing to read, because it shows exactly what was run.

Container executors: who runs the process

The NodeManager delegates starting, signalling and cleaning up processes to a container executor, set by yarn.nodemanager.container-executor.class. The default is DefaultContainerExecutor, which runs every container as the NodeManager's own Unix user. That is simple, but every application can read every other application's local files, and there is no CPU isolation.

LinuxContainerExecutor uses a setuid binary, container-executor, configured by a root-owned container-executor.cfg, to run containers as the submitting user (or a configured local user in non-secure mode). It is required for Kerberos-secured clusters, for cgroup-based CPU and memory enforcement, and for the Docker runtime. The price is setup: the binary's ownership and permissions must be exact, and banned or allowed users must be listed. See YARN and cgroups for the isolation side.

Monitoring and the kill sequence

While a container runs, the NodeManager's monitor periodically sums the memory of its process tree and compares it with the allocation. With yarn.nodemanager.pmem-check-enabled and vmem-check-enabled both true by default (virtual limit 2.1 times physical), a tree over its limit is killed; the overview explains why the virtual check is often relaxed. With cgroups and strict memory enforcement, the kernel does the policing instead.

Any kill, for memory, preemption or an AM stop request, follows the same sequence. The executor sends SIGTERM to the process group, waits yarn.nodemanager.sleep-delay-before-sigkill.ms (250 ms by default) and then sends SIGKILL. A JVM receiving SIGTERM runs its shutdown hooks, and 250 ms is not long; if your workers need to flush state on shutdown, they must be fast or you must raise the delay. A process killed this way typically shows exit code 143 (128 plus signal 15) in its own diagnostics, alongside YARN's status.

Reading exit statuses

When a container finishes, the AM receives a ContainerStatus with an exit status and a diagnostics string. Non-negative values are the process's own exit code. Negative values are ContainerExitStatus constants that mean the framework ended the container:

ValueConstantMeaning and first check
0SUCCESSExited normally
-1000INVALIDInitial value; no exit status recorded
-100ABORTEDReleased by the app or lost with its node; check node health
-101DISKS_FAILEDToo many local or log dirs bad on the node; check disks
-102PREEMPTEDTaken back for another queue; see preemption settings
-103KILLED_EXCEEDED_VMEMOver the virtual memory limit; often a false positive
-104KILLED_EXCEEDED_PMEMOver the physical memory limit; raise overhead or cut heap
-105KILLED_BY_APPMASTERThe AM asked for the stop
-106KILLED_BY_RESOURCEMANAGERThe RM ended it
-107KILLED_AFTER_APP_COMPLETIONStill running when the app finished
-108KILLED_BY_CONTAINER_SCHEDULEROpportunistic container killed to make room for a guaranteed one
-109KILLED_FOR_EXCESS_LOGSProduced too much log data

Frameworks translate these differently: Spark, for example, does not count a preempted executor as a task failure, but it does count exit code 1. When an application reports a positive code such as 1, look at the container's stderr; when it reports a negative code, look at the NodeManager's log and the diagnostics string. For preemption see YARN preemption.

Logs, aggregation and cleanup

Container output goes to yarn.nodemanager.log-dirs. Log aggregation is off by default (yarn.log-aggregation-enable is false). Without it, logs stay on the node for yarn.nodemanager.log.retain-seconds (10,800 seconds, three hours) and then disappear, which is rarely what anyone wants. With it, the NodeManager uploads each application's logs to yarn.nodemanager.remote-app-log-dir (/tmp/logs by default; move it) when the application finishes, and yarn logs reads them from there.

When a container ends, the NodeManager deletes its working directory immediately, because yarn.nodemanager.delete.debug-delay-sec defaults to 0. For debugging, set it to a few minutes on a test cluster or on a single node so you can inspect the directory; do not leave it high in production, because disks fill with dead working directories.

Worked example: a container that dies before logging

A Spark job's executors all fail on three nodes out of forty, with exit code 1 and empty stdout. The AM keeps asking for replacements, and the job crawls. The steps below narrow it down.

# 1. What did the RM and AM see?
yarn application -status application_1727700000000_0042
yarn logs -applicationId application_1727700000000_0042 \
          -containerId container_e03_1727700000000_0042_01_000007

# 2. Keep the failed container's directory around for ten minutes (NodeManager setting)
#    yarn.nodemanager.delete.debug-delay-sec = 600   (default 0: deleted at once)

# 3. On the node: the generated script shows the exact env, links and command
less /data1/yarn/local/usercache/alice/appcache/application_1727700000000_0042/\
container_e03_1727700000000_0042_01_000007/launch_container.sh

# 4. NodeManager log for localization and kill decisions
grep container_e03_1727700000000_0042_01_000007 /var/log/hadoop-yarn/*nodemanager*.log

The aggregated stderr shows Error: Could not find or load main class. The job works on other nodes, so the command is not wrong. With the debug delay set on one failing node, the working directory shows the symlink for the application jar pointing at an empty file. The NodeManager log on that node shows an earlier localization of a public resource at the same path that was interrupted when a disk briefly went read-only. The disk is back, but the cached copy is corrupt. The fix is to take that local directory out of service, clear its cache and restart the NodeManager. The long-term fix is to publish jars under versioned, immutable paths, alert on disk health (a node marks a directory bad above max-disk-utilization-per-disk-percentage, 90 percent by default), and watch per-node failure rates, because a failure concentrated on a few nodes almost always has a node-level cause.

Failure modes

SymptomCauseFix
Launch fails: resource changed on src filesystemFile overwritten after submissionVersioned, immutable paths
Every app re-downloads a big jarMarked APPLICATION or PRIVATEPUBLIC, world-readable, versioned
Exit -104 on executorsHeap plus native memory over the allocationRaise overhead, lower -Xmx
Shutdown hooks never finish250 ms between SIGTERM and SIGKILLFaster shutdown, or raise the delay
No logs after a failureAggregation off, retention 3 hoursEnable log aggregation
Failures cluster on a few nodesBad disk, corrupt cache, bad local configCheck per-node rates; decommission and clean
Capacity idle, apps waitingAM not launching allocationsFix the AM; allocations expire after 10 minutes

Trade-offs

Container size is a trade-off between packing and waste. Schedulers normalize requests up, by default to a multiple of yarn.scheduler.minimum-allocation-mb (1024 MB) and never above maximum-allocation-mb (8192 MB by default), so a request for 1,100 MB costs 2,048 MB. Many small containers localize and start more often; a few large ones fragment nodes and lose more work when one fails. The Linux executor adds security and isolation at the cost of setup, and aggressive memory checks protect neighbours at the cost of false kills. Choose each deliberately and write the choice down for the next operator. For the scheduling side, see the ResourceManager and the ApplicationMaster.

What to do next

  1. Enable log aggregation and move the remote log directory off /tmp.
  2. Publish framework jars at versioned paths and mark them PUBLIC so each node downloads them once.
  3. Check every JVM command leaves 15 to 25 percent of the container for non-heap memory.
  4. Decide between the default and Linux container executors, and document why.
  5. Make worker shutdown fit within the SIGTERM-to-SIGKILL delay or raise it.
  6. Track container exit statuses per node and per status, and alert when failures cluster on a node.
  7. Keep a debugging runbook: yarn logs, the debug delay, launch_container.sh and the NodeManager log.
Key takeaway: A YARN container is an allocation of memory and vcores on one node plus the process tree a NodeManager starts and polices against it. The AM launches it with a launch context; the NodeManager localizes files into public, private or application caches, writes launch_container.sh, starts the process through the default or Linux executor, kills it with SIGTERM then SIGKILL when needed, and reports an exit status whose negative values name the framework's reason. Use versioned public resources, leave non-heap headroom, turn on log aggregation and read exit statuses per node.