skip to content

How would you design YARN queues for a cluster shared by teams with different SLAs?

level: principalimportance: should knowfreq 32%

answer

  1. not one queue per team
  2. guarantee versus ceiling is the real lever
  3. the short query must not queue behind the long one
  4. an unenforced guarantee arrives late
  5. measure pending time, then re-argue the split

basics

~20 s

Shape queues around workload classes and SLAs rather than around org charts: give SLA-bound pipelines guaranteed capacity with limited borrowing, let ad-hoc and batch work share a large elastic queue, cap concurrency and per-user limits, and enable preemption only where redoing work is cheap.

solid answer

~50 s

Start from the SLAs, not the org chart. Group workloads by what the cluster must promise them: a small number of deadline-bound production pipelines, an interactive/ad-hoc class, and a best-effort batch class. Give the production queue a guaranteed `capacity` sized to its actual peak with a modest `maximum-capacity`, and give best-effort work a small guarantee but a high ceiling so it soaks up idle capacity. Inside the interactive queue use the `fair` ordering policy so a short query is not stuck behind a long one, and set `user-limit-factor` and `maximum-applications` so one person cannot monopolise it. Enable preemption to make the production guarantee real, but exempt work that is expensive to redo. Then instrument it: pending-container time per queue, preemption rate and utilisation tell you whether the layout matches reality. Expect to revise it — a queue design is a hypothesis about demand.

code

properties · 13 lines
properties
yarn.scheduler.capacity.root.queues=prod,adhoc,batch

yarn.scheduler.capacity.root.prod.capacity=55
yarn.scheduler.capacity.root.prod.maximum-capacity=75
yarn.scheduler.capacity.root.prod.user-limit-factor=2

yarn.scheduler.capacity.root.adhoc.capacity=25
yarn.scheduler.capacity.root.adhoc.maximum-capacity=50
yarn.scheduler.capacity.root.adhoc.ordering-policy=fair
yarn.scheduler.capacity.root.adhoc.maximum-applications=150

yarn.scheduler.capacity.root.batch.capacity=20
yarn.scheduler.capacity.root.batch.maximum-capacity=100

go deeper

for a junior

Know that a YARN cluster is divided into queues and that which queue an application lands in determines how much capacity it can get and how soon.

for a middle

Explain the mechanics you would use: guaranteed capacity versus elastic ceiling, ordering policy within a queue, and per-user limits — and what each one changes for a waiting application.

for a senior

Show the operating loop: measure pending time and preemption per queue, distinguish concurrency incidents from capacity shortages, and defend the specific limits you set on an interactive queue.

for a principal

Own the tradeoff between predictability and utilisation across the organisation, make the cost of hard isolation visible to the teams requesting it, and be clear about the problems queue design cannot solve.

## Start from promises, not from teams The common failure is one queue per team. It looks fair and it fails: each team's queue mixes a deadline-bound pipeline with somebody's exploratory notebook, so the scheduler cannot tell what matters, and fifteen small guarantees fragment the cluster so nothing large fits anywhere. Design instead around the *promise* the platform makes, then place teams inside that structure. A layout that survives contact with reality typically has three classes under `root`: - **Production / SLA** — scheduled pipelines with deadlines. Guaranteed capacity sized to measured peak concurrent demand, `maximum-capacity` only modestly above it, FIFO ordering inside (deterministic, and these jobs are large). - **Interactive / ad-hoc** — notebooks, SQL, exploration. Moderate guarantee so responsiveness is protected, `fair` ordering inside so a 30-second query never waits behind a 3-hour one, tight per-user and per-application concurrency limits. - **Best-effort / batch** — backfills, retraining, experiments. Small guarantee, ceiling near 100 so it absorbs everything idle, and the first queue preempted. Sub-queues per team belong *inside* a class, where teams share that class's guarantee, rather than at the top. ## Sizing the guarantees A guarantee should reflect measured peak concurrent demand, not a team's headcount or its budget contribution. Look at the pending-container time per queue over a few weeks: a queue whose applications rarely wait is over-provisioned, one that waits constantly is under-provisioned or has a concurrency problem rather than a capacity one. Because guarantees must sum to 100 under `root`, every increase is a decision to take capacity from somebody — which is exactly the conversation the design should force into the open rather than resolving silently at 3am. ## Elasticity is the whole point of sharing If every queue's `maximum-capacity` equals its `capacity`, you have built several small static clusters that happen to share hardware, and utilisation collapses. If every ceiling is 100, utilisation is excellent but guarantees are honoured slowly, because a borrowed container is only released when it finishes. The design decision is where each queue sits on that line: high ceilings for interruptible work, low ceilings for work whose containers are long-lived and expensive to lose. ## Preemption makes a guarantee real — at a price Without preemption a "guarantee" is a promise of eventual capacity, delivered as borrowed containers drain. With an eight-hour borrower, that is an eight-hour SLA breach. The Capacity Scheduler's scheduler monitor with a proportional preemption policy reclaims a starved queue's share by killing borrowed containers. The cost is discarded work, so calibrate rather than switch it on globally: use a grace period, prefer to preempt the best-effort class, and exempt workloads whose containers hold expensive state. Treat a persistently high preemption rate as evidence the guarantees are wrong, not as steady-state operation. ## Concurrency control is as important as capacity Most "the cluster is full" incidents are concurrency incidents. Four knobs matter: `maximum-applications` per queue caps how many can be pending or running at once; `maximum-am-resource-percent` (0.1 by default) bounds how much of a queue can be spent on coordinators, which prevents the deadlock where every resource is a coordinator waiting for workers; `user-limit-factor` stops a single user from consuming the queue's elastic headroom; and `minimum-user-limit-percent` divides a queue among active users. An interactive queue without these is one runaway notebook away from an outage. ## Placement, access and accountability Queue mappings route users and groups to queues automatically so nobody has to pass a queue name, and per-queue submit and administer ACLs stop teams landing in the production queue. Make placement automatic and auditable — a manual queue name in a job config is a bug waiting to happen. Pair this with per-queue chargeback reporting: teams argue for larger guarantees indefinitely until the cost of that guarantee is visible to them. ## Also decide what YARN cannot do for you A scheduler shares CPU and memory. It does not isolate disk or network bandwidth well, so a queue can be within its memory guarantee and still degrade neighbours through shuffle traffic or local disk saturation. It does not protect the storage layer from a queue hammering it. And it does not size the cluster — if every queue is starved, no layout fixes it. Know which problems are queue-design problems and which are capacity or workload problems, because reshuffling percentages to solve a capacity shortage burns credibility. ## Treat it as a living hypothesis Publish the layout and the reasoning, instrument pending time, preemption rate and utilisation per queue, and revisit quarterly. The strongest answer here is not a specific set of percentages — it is a design that states its assumptions, measures whether they hold, and has a defined path for a team to request a change.

  • Why is one queue per team usually the wrong starting layout?
    Because a team's queue mixes workloads with completely different promises — a deadline pipeline next to an exploratory notebook — so the scheduler has no way to prioritise what matters. It also fragments the cluster into many small guarantees, which packs badly for large containers. Classify by SLA first, then put per-team sub-queues inside a class if you need accounting.
  • How do you tell whether a starved queue needs more capacity or tighter concurrency limits?
    Look at what is waiting. If a handful of large applications each wait a long time while the queue is genuinely full, it is a capacity problem. If many applications are pending, the coordinator budget is saturated, or one user holds most of the running applications, it is concurrency — and raising the guarantee would just let the same runaway consume more. Per-queue pending-container time plus the running-application count separates the two.
  • What does queue design deliberately not solve?
    Disk and network isolation — a queue within its memory guarantee can still degrade neighbours through shuffle traffic or local disk saturation — and the storage layer, which YARN does not schedule at all. It also cannot manufacture capacity: if every queue is starved simultaneously, the cluster is undersized and no percentage split fixes it. Recognising those cases prevents pointless reallocation.
  • How would you handle a team demanding a hard, non-shareable guarantee?
    Price it honestly. Setting maximum-capacity equal to capacity gives them genuine isolation but means their idle capacity is unusable by anyone, so the cluster pays for their peak continuously. Offer the alternative — a guarantee plus preemption, which delivers their share on demand while letting others use it meanwhile — and make the utilisation cost of the hard option visible in their chargeback.

saying these in an interview costs you the question

  • Creates one queue per team and calls it fairness
  • Sets every queue's maximum-capacity to its capacity, destroying utilisation
  • Assumes a guarantee is honoured instantly without preemption
  • Sets no per-user or per-application concurrency limits
  • Reallocates percentages to solve what is actually an undersized cluster

context