A busy YARN cluster can show every container slot allocated while the nodes sit partly idle. Some of that waste is inside containers, which reserve more memory and CPU than they use. A second, quieter part sits between containers: when a short task finishes, its resources stay unused until the NodeManager's next heartbeat reaches the ResourceManager, the scheduler picks a new request, the ApplicationMaster hears about the allocation on its own heartbeat and finally launches a container. For jobs made of thousands of tasks lasting a few seconds each, that round trip can be a large fraction of each task's life.

Opportunistic containers, available in Hadoop 3.x, attack the second gap. An application can ask for containers of execution type OPPORTUNISTIC, which are allocated without waiting for free guaranteed capacity, may queue at the NodeManager until resources free up, and are killed when a guaranteed container needs the room. This article explains how they work, how to enable and request them, what the NodeManager does with them, how promotion works, and the failure modes to plan for. Facts were checked against the Apache Hadoop documentation on 2026-10-03.

Two execution types

Every YARN container has an execution type. GUARANTEED is the classic kind: the scheduler (Capacity or Fair) only allocates it when the node has unallocated resources, and once running it is stopped only by preemption for queue fairness or capacity, which YARN Preemption covers. OPPORTUNISTIC containers trade that promise for speed of dispatch:

PropertyGUARANTEEDOPPORTUNISTIC
Allocated byCapacity or Fair Scheduler in the RMOpportunistic allocator in the RM, or the NodeManager's AMRMProxy
Needs free resources at allocationYesNo; it can queue at the NM
StartsImmediately on launchWhen the NM has room, FIFO among queued opportunistic containers
Can be killed forPreemption under queue policyAny guaranteed container that needs its resources
Typical useAMs, reducers, long tasks, anything statefulShort, idempotent, retry-safe tasks such as map tasks

The documentation is explicit about one limit worth stating up front: the current implementation decides on allocated resources, not measured utilization. Opportunistic containers do not overcommit a node whose guaranteed containers are under-using their reservations; that idea is tracked separately as resource overcommitment (YARN-1011) and is listed as future work. What opportunistic containers buy you today is a queue of ready work on each node, so freed resources are reused at once instead of after a scheduling round trip, plus lower allocation latency for the applications that use them.

Centralised and distributed allocation

There are two ways an opportunistic request becomes a container. In centralised mode the ResourceManager handles it: requests arrive on the normal allocate call, the scheduler path deals with guaranteed requests, and an opportunistic allocator places opportunistic ones on the least-loaded nodes, a list the RM rebuilds from NodeManager queue reports. In distributed mode, the AMRMProxy service on each NodeManager sits between the ApplicationMaster and the RM. It implements the ApplicationMaster protocol through a chain of interceptors; its DistributedScheduler interceptor allocates opportunistic requests itself, using a list of least-loaded nodes the RM pushes to it, and forwards guaranteed requests to the RM unchanged.

ApplicationMasterrequests by ExecutionTypeAMRMProxy (on NM)DistributedSchedulerResourceManagerCapacity/Fair + Opp. allocatorNodeManager Arunning G + ONodeManager BO queue: 3NodeManager CO queue: 0ContainerSchedulerkill O (-108) for new Gallocate()GUARANTEEDleast-loaded nodeslaunch O on Claunch GDistributed mode: opportunistic requests are allocated on the node by the AMRMProxy; guaranteed requests still go to the RM.
Opportunistic allocation in distributed mode. The AM talks to the AMRMProxy on its own node; guaranteed requests continue to the RM, opportunistic ones are placed directly, queued at the target NodeManager and killed by its ContainerScheduler when guaranteed work needs the room.

Distributed mode removes the RM from the opportunistic hot path, which is where the latency win is largest, at the cost of another service on every node and a client configuration change: the documentation has job-submitting clients set yarn.resourcemanager.scheduler.address to localhost:8049, the AMRMProxy port, so the AM's scheduler traffic goes through the proxy. Centralised mode is simpler and enough when the goal is keeping node queues full rather than shaving allocation latency. See YARN ResourceManager, in depth for the heartbeat path both modes are short-circuiting.

Configuration

Everything is off by default. The properties below are the documented ones with their defaults in brackets; the values shown are a reasonable starting point for a trial, not a recommendation for every cluster.

<!-- yarn-site.xml on the ResourceManager -->
<property>
  <name>yarn.resourcemanager.opportunistic-container-allocation.enabled</name>
  <value>true</value>                     <!-- [false] -->
</property>
<property>
  <name>yarn.resourcemanager.opportunistic-container-allocation.nodes-used</name>
  <value>10</value>                       <!-- [10] least-loaded nodes considered -->
</property>
<property>
  <name>yarn.resourcemanager.nm-container-queuing.sorting-nodes-interval-ms</name>
  <value>1000</value>                     <!-- [1000] how often the list is rebuilt -->
</property>

<!-- yarn-site.xml on every NodeManager -->
<property>
  <name>yarn.nodemanager.opportunistic-containers-max-queue-length</name>
  <value>10</value>                       <!-- [0] 0 means nothing can queue -->
</property>

<!-- optional: distributed scheduling through the AMRMProxy -->
<property>
  <name>yarn.nodemanager.distributed-scheduling.enabled</name>
  <value>true</value>                     <!-- [false] -->
</property>
<property>
  <name>yarn.nodemanager.amrmproxy.address</name>
  <value>0.0.0.0:8049</value>             <!-- [0.0.0.0:8049] -->
</property>

The queue length default of zero is the trap: with the RM side enabled and the NM side left alone, opportunistic containers cannot queue anywhere, which defeats the design. Three further RM properties control load shedding of those queues: yarn.resourcemanager.nm-container-queuing.min-queue-length [5], ...max-queue-length [15] and ...queue-limit-stdev [1.0], which together bound how long any node's queue may grow relative to the cluster mean. Restart the RM and NMs after changing them, and confirm on the RM web UI or REST API that the setting took effect before benchmarking.

Requesting opportunistic containers

Applications choose the type per request through an ExecutionTypeRequest, which carries the type and an enforce flag. With enforcement off, the scheduler may satisfy the request with either type; with it on, only the requested type will do. MapReduce exposes this as a single job property, mapreduce.job.num-opportunistic-maps-percent, the percentage of map tasks to run as opportunistic. Reducers stay guaranteed, which is sensible because a killed reducer loses all its shuffle progress.

// Inside a custom ApplicationMaster using AMRMClient
Resource cap = Resource.newInstance(2048, 1);            // 2 GB, 1 vcore
Priority pri = Priority.newInstance(10);
ExecutionTypeRequest opp =
    ExecutionTypeRequest.newInstance(ExecutionType.OPPORTUNISTIC, true);

AMRMClient.ContainerRequest req = new AMRMClient.ContainerRequest(
    cap, null, null, pri, true, null, opp);  // nodes, racks, priority, relaxLocality, label expr
amRMClient.addContainerRequest(req);

// Later, in the allocate loop
for (Container c : response.getAllocatedContainers()) {
    if (c.getExecutionType() == ExecutionType.OPPORTUNISTIC) {
        launchShortTask(c);          // idempotent work only
    } else {
        launchCriticalTask(c);
    }
}

Check the constructor overloads in the AMRMClient.ContainerRequest javadoc for your Hadoop version; the shape has varied across 3.x releases, and some versions also offer a builder. The distributed shell example accepts -container_type OPPORTUNISTIC and is the quickest way to watch the feature work before writing any code.

What the NodeManager does with them

On the NodeManager, the ContainerScheduler decides what runs. When an opportunistic container arrives and the node's allocated resources leave room, it starts at once. Otherwise it waits in a FIFO queue bounded by the max queue length; how a container beyond that bound is handled has varied between releases, so test it on yours rather than assuming it waits. When a guaranteed container arrives and there is not enough room, the scheduler kills running opportunistic containers until it fits, and those containers finish with exit status -108, KILLED_BY_CONTAINER_SCHEDULER. The table of exit statuses in YARN Containers, in depth places -108 among the other kill reasons.

This makes the guarantee concrete: guaranteed containers never wait for opportunistic ones. It also means opportunistic work has no floor. On a node where guaranteed containers keep arriving, an opportunistic container can be started and killed repeatedly, wasting the work it did each time. The documentation lists smarter queue reordering (YARN-5886), out-of-order killing (YARN-5887) and pausing instead of killing (YARN-5292) as future work, so do not assume any of them on your release without checking its release notes. The NodeManager behaviours around health and restarts that interact with this are in YARN NodeManager, in depth.

Promotion and demotion

An application can ask the RM to change a running container's type. ContainerUpdateType has four values: INCREASE_RESOURCE, DECREASE_RESOURCE, PROMOTE_EXECUTION_TYPE and DEMOTE_EXECUTION_TYPE. A promotion request names the container, its current version and the target type:

UpdateContainerRequest up = UpdateContainerRequest.newInstance(
    c.getVersion(), c.getId(),
    ContainerUpdateType.PROMOTE_EXECUTION_TYPE,
    null,                                   // no resource change
    ExecutionType.GUARANTEED);
amRMClient.requestContainerUpdate(c, up);
// The result arrives on a later allocate() response as an updated container.

Promotion only succeeds when the queue has guaranteed headroom, so it is a request, not a command. A useful pattern is to start a task opportunistically for low latency, then promote it when it turns out to run long; the distributed shell's -promote_opportunistic_after_start flag exercises this path. Older documentation pages list promotion (YARN-5085) under future work while the API exists in current javadoc, so test it on your exact version.

Worked example: a short-map MapReduce job

Consider a 50-node cluster running a nightly MapReduce job with 20,000 map tasks of about 8 seconds each, alongside interactive guaranteed work. With guaranteed containers only, each map slot sits empty from task end until the next allocation is launched. Measure that gap from the job history: compare the sum of task durations to the slot-seconds the job held. Suppose it averages 1.5 seconds per task, about 16 percent of an 8-second task.

  1. Enable centralised allocation on the RM and set the NM queue length to 4, so each node holds a few ready maps.
  2. Run the job with mapreduce.job.num-opportunistic-maps-percent=50. Half the maps are opportunistic; the other half keep a guaranteed floor of progress.
  3. Compare wall-clock time and the count of attempts killed with -108. If kills exceed a few percent of opportunistic attempts, guaranteed load is too heavy for this share; reduce the percentage.
  4. Raise to 80 percent only when the kill rate stays low, and keep reducers guaranteed.

The expected win is bounded by the gap you measured: if the inter-task gap is 16 percent, opportunistic queuing cannot save more than roughly that, and kills eat into it. If the gap is small because tasks are long, opportunistic containers buy little and add risk.

Failure modes

  • Kills counted as failures. A framework that treats -108 as a task failure can fail a job after a few kills. MAPREDUCE-7205 changed MapReduce to treat it as a kill; check that your version includes the change, and check how any other framework you use maps -108.
  • Feature enabled but nothing queues. The NM max queue length still defaults to 0. Symptoms are opportunistic requests that never start or are rejected.
  • Non-idempotent opportunistic work. A task that writes to an external system without idempotence repeats its side effect after each kill and retry.
  • Starvation under steady guaranteed load. Opportunistic containers are killed repeatedly and the job runs slower than with no opportunistic share at all.
  • Expecting overcommit. Teams enable the feature to use idle memory inside guaranteed reservations, see no utilization change, and conclude it is broken. It allocates on allocated resources; overcommit is separate work.
  • Client configuration left unchanged. In distributed mode, clients that still submit with the normal RM scheduler address send AM traffic around the AMRMProxy, so distributed allocation never happens.

Operating it

Treat the opportunistic share as a tunable with a feedback signal. The signal is the ratio of opportunistic attempts killed with -108 to opportunistic attempts started, per job and per queue. Collect it from job history or from container exit statuses in the NodeManager logs, and graph it next to cluster allocated memory. A kill ratio near zero with long NM queues means you can raise the share; a rising ratio during business hours means guaranteed load is pushing opportunistic work out, and the share should drop or the jobs should move to quieter hours.

Watch queue lengths too. Persistently full queues on a few nodes while others are empty suggest the least-loaded list is stale; shortening the sorting interval or lowering the max queue length spreads work more evenly. Roll out per job, not cluster-wide: start with one batch pipeline whose tasks are known to be idempotent, keep a guaranteed-only run as a baseline, and compare wall-clock time across a week of normal load before widening the change.

Trade-offs

Opportunistic containers trade predictability for throughput and latency on short, retry-safe work. They help most when tasks are short, numerous and idempotent, and when the inter-task gap is measurable. They help least for long tasks, stateful work and clusters whose guaranteed load leaves no room. Distributed scheduling adds latency savings at the cost of an extra service per node. Keep the AM itself, reducers, and anything holding exclusive external state on guaranteed containers.

What to do next

  1. Measure the inter-task gap of your largest short-task jobs from job history before enabling anything.
  2. Enable RM-side allocation and set a small non-zero NM queue length on a test cluster; confirm the settings through the RM.
  3. Run the distributed shell with -container_type OPPORTUNISTIC and watch containers queue and start.
  4. Confirm your MapReduce version treats -108 as a kill, then trial mapreduce.job.num-opportunistic-maps-percent at 25 to 50.
  5. Track kill counts with exit status -108 per job and alert when they climb.
  6. Try distributed scheduling only if allocation latency, not queue depth, is the remaining bottleneck.
  7. Read the Hadoop 3 features overview for the other changes that shipped alongside this one.
Key takeaway: Opportunistic containers keep a queue of ready work on each NodeManager so freed resources are reused without a scheduling round trip. They are allocated on allocated resources, not measured utilization, and are killed with exit status -108 when guaranteed work needs room. Enable both the RM allocator and a non-zero NM queue length, use them only for short idempotent tasks, confirm your framework counts -108 as a kill, and size the opportunistic share from the inter-task gap you measured.