The YARN Fair Scheduler shares a Hadoop cluster between queues so that, over time, each gets a share proportional to its weight, while idle capacity is lent to whoever can use it. Most teams enable it with defaults, then spend months fighting symptoms: a queue that never reaches its share, applications stuck in ACCEPTED, containers killed at random, or jobs landing in queues nobody created.

This page is the operator's view. The Fair Scheduler architecture page explains lending, dominant resource fairness and the life of a request. Here we build a real allocation file, do the fair-share and preemption arithmetic by hand, and map each common symptom to the setting that causes it. Property and element names follow the Apache Hadoop documentation for the Fair Scheduler.

Advertisement

How the scheduler decides, in one paragraph

Every update interval (yarn.scheduler.fair.update-interval-ms, 500 ms by default) the scheduler recomputes a fair share for every queue in the tree. When a NodeManager heartbeats with free capacity, the scheduler walks down from root, at each level choosing the child that is furthest below its fair share according to that parent's scheduling policy, until it reaches an application with a pending request that fits. It checks queue limits, application-master limits and locality, then allocates a container. If preemption is enabled, a separate process finds queues that have been starved for too long and reclaims containers from queues above their share. Almost every tuning question is about one of those four steps: placement, share computation, assignment, or preemption.

From submission to container: placement, fair-share computation, assignment, preemptionClient submitqueue=? user, groupsPlacement rulesspecified, primaryGroup, defaultLeaf queueroot.etl / ml / adhocrejectif chain ends in rejectmatchFair-share computation (every update interval)find ratio r so shares sum to cluster sizeshare = weight x r, raised to minResources,capped at maxResources and at demandinstantaneous: active queues onlysteady: all configured queuesAssignment on NodeManager heartbeatwalk tree, pick queue furthest below sharethen app by schedulingPolicy (fair, drf, fifo)check maxAMShare, maxRunningApps, localityassignmultiple / max.assign: containers per beatPreemption (off by default: yarn.scheduler.fair.preemption=true)only when cluster utilisation exceeds 0.8 (cluster-utilization-threshold)starved: usage below fairSharePreemptionThreshold x fair share for fairSharePreemptionTimeout,or below minResources for minSharePreemptionTimeoutvictims from queues over their share whose allowPreemptionFrom is truefreed containers
The four stages an operator configures. Placement decides the queue, the share computation decides the target, assignment hands out containers, preemption corrects long-lived imbalances.

Fair share arithmetic, worked

Take a cluster of 1,200 GB of memory and 300 vcores with three top-level queues: etl with weight 3, ml with weight 2 and adhoc with weight 1. The scheduler finds a ratio r such that each queue's share is its weight times r, raised to its minResources, capped at its maxResources and at what it is actually asking for, and the shares sum to the cluster size. With no minimums, maximums or demand caps binding, that is plain proportional division.

Situationetl (w=3)ml (w=2)adhoc (w=1)
All three active, all hungry600 GB (50%)400 GB (33%)200 GB (17%)
Only ml and adhoc active0800 GB (67%)400 GB (33%)
All active; adhoc needs only 60 GB684 GB456 GB60 GB
Steady fair share (always)600 GB400 GB200 GB

The second and third rows are the instantaneous fair share: only queues with running applications take part, and a queue asking for less than its proportional slice releases the rest to the others in proportion to their weights (1,140 GB split 3:2). The last row is the steady fair share, computed over all configured queues whether active or not. The web UI shows both; scheduling and preemption use the instantaneous value. A frequent support question is why adhoc runs at 30 percent of the cluster when its configured share is 17 percent. The 17 percent is its steady share; with etl idle its instantaneous share is 33 percent, so at 30 percent it is still below its fair share in the sense scheduling and preemption use. It is using loaned capacity legitimately, and that loan is reclaimed only when etl becomes active again.

Now add minResources of 200 GB to etl. When etl is active its share cannot fall below 200 GB even if weights would give it less, which matters when many queues are active. Add maxResources of 800 GB to ml, and in the second row ml is capped at 800 GB anyway; if adhoc were idle too, the remaining 400 GB would simply stay unused by ml. The maxResources cap is absolute, not a share.

Advertisement

A worked allocation file

The allocation file (fair-scheduler.xml by default, set by yarn.scheduler.fair.allocation.file) holds queues, limits, defaults and placement. The ResourceManager reloads it within about 10 to 15 seconds of a change, so most tuning needs no restart. Resources should use the recommended memory-mb=..., vcores=... form; in minResources an unspecified resource defaults to zero and in maxResources to unlimited.

<?xml version="1.0"?>
<allocations>
  <defaultQueueSchedulingPolicy>fair</defaultQueueSchedulingPolicy>
  <defaultFairSharePreemptionThreshold>0.5</defaultFairSharePreemptionThreshold>
  <defaultFairSharePreemptionTimeout>300</defaultFairSharePreemptionTimeout>
  <queueMaxAMShareDefault>0.3</queueMaxAMShareDefault>
  <userMaxAppsDefault>10</userMaxAppsDefault>

  <queue name="etl">
    <weight>3</weight>
    <minResources>memory-mb=204800, vcores=48</minResources>
    <minSharePreemptionTimeout>120</minSharePreemptionTimeout>
    <schedulingPolicy>drf</schedulingPolicy>
    <allowPreemptionFrom>false</allowPreemptionFrom>
    <aclSubmitApps>etl-svc etl-team</aclSubmitApps>
  </queue>

  <queue name="ml">
    <weight>2</weight>
    <maxResources>memory-mb=819200, vcores=200</maxResources>
    <schedulingPolicy>drf</schedulingPolicy>
  </queue>

  <queue name="adhoc">
    <weight>1</weight>
    <maxRunningApps>40</maxRunningApps>
    <maxAMShare>0.2</maxAMShare>
    <fairSharePreemptionThreshold>0.8</fairSharePreemptionThreshold>
  </queue>

  <queuePlacementPolicy>
    <rule name="specified" create="false"/>
    <rule name="primaryGroup" create="false"/>
    <rule name="default" queue="adhoc"/>
  </queuePlacementPolicy>
</allocations>

Read it line by line. The defaults set every queue's fair-share preemption threshold to 0.5 with a 300-second timeout and limit application masters to 30 percent of a queue's share. etl is the production queue: the largest weight, a guaranteed minimum it can reclaim within two minutes, DRF so its memory-heavy and CPU-heavy jobs are compared by dominant resource, and allowPreemptionFrom false so its containers are never victims. ml is capped so a big training sweep cannot absorb the cluster. adhoc limits concurrent applications and application masters, and uses a stricter 0.8 threshold so interactive users get their share back quickly.

The ACL deserves a note: queue ACLs are inherited downward and root allows everyone unless you restrict it, so a submit ACL on a leaf does nothing if root still grants *. Set root's ACL to a single space, meaning nobody, and grant access on the leaves.

Placement rules: which queue does a job land in?

The queuePlacementPolicy is an ordered chain; the first rule that yields a queue wins. In the example, specified honours an explicit queue from the client but, with create="false", only if it already exists; a request for default falls through. primaryGroup places the job into a queue named after the user's primary group if one exists, so members of the ml group land in ml without typing it. Everything else lands in adhoc through the final default rule; ending with reject instead turns a missing queue into a submission error, which is stricter but louder.

If you want each ad hoc user to get an equal slice rather than letting one user with 30 jobs dominate, the nestedUserQueue rule creates a queue named after the user, such as root.adhoc.alice, under the queue its nested rule returns. It applies only if that nested rule returns a parent queue, meaning one declared with type="parent" or one that already has a child. Because maxAMShare can only be set on leaf queues, converting adhoc to a parent means the generated per-user leaves take queueMaxAMShareDefault instead, so review both together.

Two settings interact with this chain. yarn.scheduler.fair.allow-undeclared-pools (default true) lets rules create queues on the fly, and yarn.scheduler.fair.user-as-default-queue (default true) makes the username the default queue when no policy is configured. Leaving both at their defaults with no placement policy is how clusters end up with hundreds of top-level queues named after users, each taking an equal share. Periods in usernames are replaced with _dot_ in generated queue names.

Preemption, with numbers

Preemption is off unless yarn.scheduler.fair.preemption is true, and even then it acts only while overall cluster utilisation is above yarn.scheduler.fair.preemption.cluster-utilization-threshold, 0.8 by default; on a lightly loaded cluster, waiting is cheaper than killing. A queue is starved for fair share when its usage stays below fairSharePreemptionThreshold times its instantaneous fair share for fairSharePreemptionTimeout seconds, and starved for min share when it stays below minResources for minSharePreemptionTimeout seconds.

Work it through. With all three queues active, adhoc's fair share is 200 GB. Its threshold is 0.8, so it is starved while below 160 GB. Suppose ml borrowed the whole cluster overnight and adhoc users arrive at 9:00 with plenty of demand. Containers finish naturally at some rate; if adhoc is still below 160 GB after 300 seconds, the scheduler selects containers from queues above their fair share (ml, and not etl, which is protected) and asks for them back, killing them after a grace period if the application does not release them. Meanwhile etl, if it arrived at the same time and sat below its 200 GB minimum, would trigger min-share preemption after only 120 seconds.

The tuning trade-off is between latency for the starved queue and wasted work for the victim. A killed map task costs its runtime; a killed Spark executor can cost a stage retry and its cached data. Short timeouts plus a high threshold give interactive users fast recovery at the price of more kills, so reserve them for queues whose jobs are short. The YARN preemption guide covers the kill mechanism and how applications should handle preemption notices.

Application masters, running-app limits and the ACCEPTED trap

Every YARN application starts with an application master container. maxAMShare (default 0.5 per queue) limits the fraction of a queue's fair share that application masters may occupy. The limit protects against a pathological state in which a queue is full of application masters waiting for task containers that can never be allocated because the masters used up the share.

The flip side is a common stall. adhoc's instantaneous share may be small when the cluster is busy: with 60 GB of share and maxAMShare of 0.2, only 12 GB can hold application masters. If each Spark driver asks for 4 GB, the fourth application sits in ACCEPTED even though the cluster has free memory elsewhere, because the share, not the cluster, is the denominator. The fix is to lower driver memory, raise the queue's weight or minimum, or raise maxAMShare for that queue; maxRunningApps gives a clearer, count-based limit for the same problem.

Assignment speed and locality

By default one container is assigned per node heartbeat. On clusters with large nodes and short tasks this is a throughput limit, so yarn.scheduler.fair.assignmultiple allows several per heartbeat, with yarn.scheduler.fair.dynamic.max.assign (default true) choosing how many from the node's free capacity or yarn.scheduler.fair.max.assign setting a fixed cap. Assigning many at once packs work onto the first nodes to heartbeat, which helps utilisation but can create hot nodes.

Locality delays are set by yarn.scheduler.fair.locality.threshold.node and .rack, expressed as a fraction of cluster nodes whose scheduling opportunities an application passes up while waiting for a better placement; the default of -1.0 means no waiting. For HDFS-heavy batch jobs, a small value such as 0.1 to 0.3 can raise data-local tasks; for workloads that read from object storage it buys nothing. yarn.scheduler.fair.sizebasedweight weights applications by the logarithm of their memory demand, favouring large jobs; it is off by default and rarely what people want.

Observing the scheduler

The ResourceManager UI's scheduler page shows, for every queue, used resources, instantaneous and steady fair share, minimum and maximum. The same data is available as JSON from the REST endpoint, which is what dashboards and alerts should read:

# Queue tree with fair shares and usage
curl -s http://rm-host:8088/ws/v1/cluster/scheduler | jq '.scheduler.schedulerInfo'

# Move a misplaced application without restarting it
yarn application -movetoqueue application_1727650000000_0042 -queue root.adhoc

Alert on three things: a queue whose usage stays under half its instantaneous fair share for longer than its preemption timeout while it has pending demand (preemption is off or misconfigured), applications in ACCEPTED for more than a few minutes (AM share or running-app limits), and container preemption counts per queue (a rising number means timeouts are too aggressive). The ResourceManager page describes where these metrics come from.

Failure modes

  • Preemption silently off. Timeouts and thresholds are set but yarn.scheduler.fair.preemption is false, so starved queues wait for containers to finish naturally, sometimes for hours.
  • Queue explosion. Default placement with undeclared pools creates a queue per user, flattening every weight you designed.
  • Stuck in ACCEPTED. Small instantaneous shares plus maxAMShare block new application masters while the cluster has free capacity.
  • Memory-only fairness. The default fair policy compares memory only; CPU-heavy jobs can saturate vcores unnoticed. Use drf where CPU matters.
  • Preemption churn. Short timeouts on queues with long tasks kill hours of work repeatedly. Protect such queues or lengthen their timeouts.
  • Unreadable allocation file. A malformed edit is rejected on reload and the old configuration stays in force; the change you think you made never happened. Check the ResourceManager log after every edit.

Fair or Capacity, and migrating

The Capacity Scheduler expresses the same ideas as percentages of the parent with elastic maxima, and has become the default in many distributions. The Fair Scheduler's weights, automatic per-user queues and simple configuration still suit clusters shared by many teams with bursty demand. If you are moving, the yarn fs2cs converter documented by Cloudera reads yarn-site.xml and fair-scheduler.xml and emits a capacity-scheduler.xml and updated yarn-site.xml; check that your Hadoop distribution includes it and treat its output as a draft. Compare with the Capacity Scheduler guide before deciding.

What to do next

  1. Pull /ws/v1/cluster/scheduler and list every queue with weight, min, max, instantaneous and steady fair share; delete queues nobody owns.
  2. Write an explicit queuePlacementPolicy ending in reject or a known default, and restrict the root ACL.
  3. Switch queues running Spark or mixed workloads to drf.
  4. Enable preemption deliberately: set thresholds and timeouts per queue, protect production queues with allowPreemptionFrom false, and watch kill counts for a week.
  5. Find applications that sat in ACCEPTED and check them against maxAMShare and maxRunningApps.
  6. Add alerts for starved queues, ACCEPTED duration and preemption rate, and re-run the fair-share arithmetic whenever you add a queue.
Key takeaway: The Fair Scheduler divides a cluster by weight among active queues, lends idle capacity, and uses preemption to claw it back. Instantaneous fair share, computed over active queues and capped by demand, drives scheduling; steady fair share is for planning. Make placement explicit, use DRF where CPU matters, enable preemption on purpose with per-queue thresholds and timeouts and protected production queues, and watch application-master limits, which are the usual cause of jobs stuck in ACCEPTED on an apparently free cluster. Edit the allocation file, check the log, and verify through the REST API.