skip to content

How do YARN's Capacity Scheduler and Fair Scheduler differ when sharing one cluster?

level: middleimportance: should knowfreq 45%

answer

  1. one divides in advance, one divides now
  2. guarantee plus a borrowing ceiling
  3. the small query behind the huge job
  4. inside a queue it is FIFO by default
  5. reclaiming a share costs redone work

basics

~20 s

Both divide a YARN cluster into hierarchical queues, but the Capacity Scheduler starts from guaranteed percentages per queue with an elastic ceiling, while the Fair Scheduler starts from equal sharing among running applications and pulls resources back toward each queue's fair share over time.

solid answer

~50 s

The scheduler is pluggable through `yarn.resourcemanager.scheduler.class`; the **Capacity Scheduler** is the default in Apache Hadoop 3 (FIFO also exists but is single-tenant in practice). Capacity thinks in **guarantees**: each queue is configured with a `capacity` percentage it is entitled to, plus a `maximum-capacity` ceiling it may elastically borrow up to when neighbours are idle, and inside a queue applications are FIFO by default. **Fair** thinks in **shares**: it hands resources to whichever running application is furthest below its fair share, so a small job submitted behind a large one starts making progress almost immediately without waiting for the big one to finish. Both now support hierarchical queues, weights, dominant-resource fairness and preemption, so they have largely converged — the practical difference is that Capacity is the actively developed one, and vendors have standardised on it with migration tooling for existing Fair configurations.

code

properties · 9 lines
properties
yarn.resourcemanager.scheduler.class=\
  org.apache.hadoop.yarn.server.resourcemanager.scheduler.capacity.CapacityScheduler

yarn.scheduler.capacity.root.queues=etl,adhoc
yarn.scheduler.capacity.root.etl.capacity=70
yarn.scheduler.capacity.root.etl.maximum-capacity=100
yarn.scheduler.capacity.root.adhoc.capacity=30
yarn.scheduler.capacity.root.adhoc.maximum-capacity=60
yarn.scheduler.capacity.root.adhoc.ordering-policy=fair

go deeper

for a junior

Know that YARN divides a cluster into queues and that the scheduler decides which queue's application gets the next container; recognise the two named schedulers.

for a middle

Explain the two mental models — pre-agreed guarantees with elastic borrowing versus equalising among what is running now — and name the capacity and maximum-capacity settings that express the first.

for a senior

Demonstrate the operational consequences: why a guarantee without preemption arrives late, what preemption costs in redone work, and how within-queue ordering policy changes interactive responsiveness.

for a principal

Own the allocation policy for the organisation: how queue shares map to funding or SLAs, how much elasticity to permit, and whether the cluster is being run for predictability or for utilisation.

## The scheduler is a plugin The ResourceManager delegates every allocation decision to a scheduler chosen by `yarn.resourcemanager.scheduler.class`. Three ship with Hadoop: - **FIFO Scheduler** — one queue, first-come-first-served. One large job blocks everything behind it. Useful only for a single-user cluster or a test rig. - **Capacity Scheduler** — the default in Apache Hadoop 3, configured in `capacity-scheduler.xml`. - **Fair Scheduler** — configured through a separate allocation file, commonly `fair-scheduler.xml`. ## Capacity Scheduler: guarantees first The mental model is that the cluster is *divided among organisations in advance*. Queues form a tree under `root`, and each is configured with the percentage of its parent it is entitled to: - `yarn.scheduler.capacity.root.<q>.capacity` — the guaranteed share. Siblings must sum to 100. - `yarn.scheduler.capacity.root.<q>.maximum-capacity` — the elastic ceiling. A queue may borrow idle capacity from siblings up to this number, and is squeezed back toward its guarantee as they become active. - `user-limit-factor` (default 1) — how far one user may exceed the queue's configured capacity, which is why a single user often cannot use the elastic headroom. - `maximum-am-resource-percent` (default 0.1) — the coordinator budget described above. - `ordering-policy` — `fifo` by default within a queue, but it can be set to `fair` so applications inside one queue share rather than queue up. The promise is *predictability*: a team that owns 30% knows it can always get 30%, and knows it will get more when the cluster is quiet. ## Fair Scheduler: shares over time The mental model is that *everything running right now should get an equal slice*. On each scheduling opportunity the Fair Scheduler gives the container to whichever queue or application is furthest below its computed fair share. With one job running it gets the whole cluster; when a second starts, resources drain toward it as the first job's containers finish (and immediately, if preemption is enabled). Queues have weights rather than fixed percentages, plus optional `minResources` and `maxResources` bounds, and each queue chooses a policy — FIFO, fair on memory, or `drf` for dominant resource fairness across memory and CPU. Its historical selling point was interactive responsiveness: an analyst's five-minute query submitted behind an eight-hour ETL job starts progressing at once instead of waiting. ## Where they have converged Older comparisons overstate the gap. Both schedulers today offer hierarchical queues, per-queue ACLs, placement rules that map users to queues, dominant-resource fairness across memory and vcores, and preemption. The Capacity Scheduler can use `DominantResourceCalculator` instead of the memory-only `DefaultResourceCalculator` to schedule on CPU as well as memory, and its `fair` ordering policy inside a queue reproduces much of the Fair Scheduler's within-queue behaviour. In practice, Hadoop distributions consolidated on the Capacity Scheduler and provided a conversion tool for existing Fair configurations, so new clusters overwhelmingly run Capacity and the Fair Scheduler is legacy knowledge. ## Preemption: the tradeoff neither avoids Without preemption, a queue's "guarantee" is only honoured as containers naturally finish. If a neighbour borrowed capacity for an eight-hour job, your guaranteed share arrives eight hours late. Preemption fixes that by killing borrowed containers to reclaim a starved queue's entitlement — the Capacity Scheduler does this with a scheduler monitor running a proportional preemption policy. The cost is real: a preempted container's work is lost and must be redone, so aggressive preemption converts a fairness problem into a wasted-compute problem. Sensible practice is to enable it with a grace period, exempt queues whose work is expensive to redo, and treat a high preemption rate as a signal that the queue layout is wrong rather than as normal operation. ## How to answer "which would you pick" Say what the workload needs. If teams need contractual, predictable shares and chargeback, guarantee-first allocation with modest elasticity fits. If the cluster is dominated by many interactive queries of wildly different sizes, share-based allocation or the `fair` ordering policy inside a queue is what keeps small queries responsive. Then note that on Hadoop 3 you would implement either with the Capacity Scheduler, because it is the default and the maintained one — the choice today is about queue design and ordering policy, not about installing a different scheduler.

  • What is the difference between a queue's capacity and its maximum-capacity?
    Capacity is the guarantee — the share the queue can always obtain when it asks. Maximum-capacity is the elastic ceiling it may borrow up to while siblings are idle, and it is squeezed back toward the guarantee as they become active. Setting maximum-capacity equal to capacity gives hard isolation with no borrowing; setting it to 100 maximises utilisation at the cost of neighbours waiting for containers to free up.
  • Why would you enable preemption, and what does it cost?
    Without it, a starved queue only regains its guaranteed share as borrowed containers finish naturally, so a long-running neighbour can delay a guarantee by hours. Preemption kills borrowed containers to reclaim capacity immediately. The cost is that the killed container's work is discarded and re-run, so it converts a latency problem into wasted compute — use a grace period, and treat a high preemption rate as evidence that the queue layout is wrong.
  • What does switching the Capacity Scheduler to DominantResourceCalculator change?
    The default resource calculator considers memory only, so CPU-heavy applications are scheduled as if their vcore demand were free and a node can be oversubscribed on CPU. DominantResourceCalculator applies dominant resource fairness across memory and vcores, comparing each application on whichever resource it consumes the largest share of. It matters on mixed clusters where some workloads are CPU-bound and others memory-bound.

saying these in an interview costs you the question

  • Says the Fair Scheduler is the Hadoop 3 default
  • Believes only one of the two supports hierarchical queues
  • Treats a queue's guaranteed capacity as a hard cap it can never exceed
  • Thinks preemption is free and loses no work
  • Claims applications inside a Capacity Scheduler queue always share fairly

context