A YARN cluster shared by several teams needs a way to say who gets what. The Capacity Scheduler, the default scheduler in Hadoop 3, answers with a tree of queues. Each queue gets a guaranteed share of its parent, can borrow idle capacity up to a ceiling, and applies limits to the users and applications inside it. The tree is where the organisation's priorities meet the scheduler, and most multi-tenant problems, such as a team that can never get resources, a queue that everyone can submit to, or jobs stuck in ACCEPTED, come from how the tree was designed.
This page is about designing and running that tree. It works through the capacity arithmetic on a concrete cluster, shows the configuration, explains the rules that surprise people, particularly ACL inheritance and application-master limits, and covers changing the tree without a restart. The scheduler's internals are covered in the Capacity Scheduler architecture page.
Parents divide, leaves run
Every hierarchy starts at root. A queue with children is a parent queue; a queue without children is a leaf. Applications can only be submitted to leaf queues. Parent queues never hold applications: they exist to divide capacity among their children and to decide how idle capacity is lent.
Queues are named by their path, such as root.prod.etl. A leaf name only has to be unique within its parent, so root.prod.adhoc and root.analytics.adhoc can both exist. That is convenient, but it makes short names ambiguous, so use full paths in placement rules and scripts.
The hierarchy expresses two kinds of decision. The split at each level says how resources are shared when everyone is busy. The ceilings, user limits and ACLs say what any one tenant can do when others are idle. Keep the two in mind separately: most mistakes come from tuning one while thinking about the other.
Capacity arithmetic, worked
Take a cluster with 1,000 GB of memory and 500 vcores available to YARN. In percentage mode, each queue's capacity is a percentage of its parent, and siblings must sum to 100. The absolute guarantee of a leaf is the product along its path. root.prod at 70% gets 700 GB. root.prod.etl at 60% of prod gets 0.7 × 0.6 = 42% of the cluster, 420 GB and 210 vcores. root.analytics.bi at 50% of analytics gets 0.25 × 0.5 = 12.5%, or 125 GB.
The guarantee is what the queue gets under contention. When siblings are idle, a queue can grow beyond it, up to maximum-capacity. The default of -1 means 100%. A child's maximum is also bounded by its parent's maximum, so in the example, where root.analytics has a maximum of 50%, adhoc can never exceed 500 GB however idle the cluster is, even though its own maximum is the default.
The Capacity Scheduler supports two other ways to express capacity. Absolute mode gives a queue an explicit amount, such as [memory=102400,vcores=50], which suits contracts expressed in hardware. Weight mode gives a relative weight, such as 2.0w, and shares the parent in proportion to sibling weights, so adding a sibling does not force you to rebalance numbers that must sum to 100. Children of one parent must use the same mode. Weights are the easiest to maintain for trees that change often.
The configuration
The same tree, with the limits and ACLs discussed below, in capacity-scheduler.xml:
<!-- capacity-scheduler.xml (excerpt) -->
<property><name>yarn.scheduler.capacity.root.queues</name><value>prod,analytics,default</value></property>
<property><name>yarn.scheduler.capacity.root.prod.capacity</name><value>70</value></property>
<property><name>yarn.scheduler.capacity.root.analytics.capacity</name><value>25</value></property>
<property><name>yarn.scheduler.capacity.root.default.capacity</name><value>5</value></property>
<property><name>yarn.scheduler.capacity.root.prod.queues</name><value>etl,serving</value></property>
<property><name>yarn.scheduler.capacity.root.prod.etl.capacity</name><value>60</value></property>
<property><name>yarn.scheduler.capacity.root.prod.serving.capacity</name><value>40</value></property>
<property><name>yarn.scheduler.capacity.root.prod.serving.maximum-capacity</name><value>60</value></property>
<property><name>yarn.scheduler.capacity.root.analytics.queues</name><value>adhoc,bi</value></property>
<property><name>yarn.scheduler.capacity.root.analytics.adhoc.capacity</name><value>50</value></property>
<property><name>yarn.scheduler.capacity.root.analytics.bi.capacity</name><value>50</value></property>
<property><name>yarn.scheduler.capacity.root.analytics.maximum-capacity</name><value>50</value></property>
<property><name>yarn.scheduler.capacity.root.analytics.adhoc.user-limit-factor</name><value>2</value></property>
<property><name>yarn.scheduler.capacity.root.analytics.adhoc.minimum-user-limit-percent</name><value>25</value></property>
<property><name>yarn.scheduler.capacity.root.analytics.bi.maximum-am-resource-percent</name><value>0.3</value></property>
<!-- ACLs: close root first, then grant per branch -->
<property><name>yarn.scheduler.capacity.root.acl_submit_applications</name><value> </value></property>
<property><name>yarn.scheduler.capacity.root.acl_administer_queue</name><value>yarnadmin</value></property>
<property><name>yarn.scheduler.capacity.root.prod.acl_submit_applications</name><value>etlsvc,servesvc prod-eng</value></property>
<property><name>yarn.scheduler.capacity.root.analytics.acl_submit_applications</name><value> analysts</value></property>
<property><name>yarn.scheduler.capacity.root.default.acl_submit_applications</name><value>*</value></property>
<!-- Placement -->
<property><name>yarn.scheduler.capacity.queue-mappings</name>
<value>u:etlsvc:root.prod.etl,u:servesvc:root.prod.serving,g:analysts:root.analytics.adhoc</value></property>
<property><name>yarn.scheduler.capacity.queue-mappings-override.enable</name><value>false</value></property>Every property follows yarn.scheduler.capacity.<queue-path>.<property>. Note serving has a maximum of 60%, applied to prod's own ceiling of 100% of the cluster, so 600 GB: a latency-sensitive leaf is often capped so it cannot absorb capacity that etl will need back through preemption, while analytics as a whole is capped at 50% so ad hoc work can never take most of the cluster.
User limits inside a leaf
Within a leaf, the scheduler also limits what one user can take. user-limit-factor is a multiple of the queue's capacity that a single user may use. The default is 1, which means one user can use at most the queue's guaranteed capacity, even when the queue itself could grow into idle cluster capacity. This is the most common reason a lone user's job "won't use the empty cluster". In the example, adhoc sets it to 2 so a single analyst can use up to 250 GB when the cluster is quiet.
minimum-user-limit-percent works the other way: it sets the minimum share each active user is guaranteed when several compete. At 25, the first four active users can each be held to at least a quarter of the queue; a fifth user waits for capacity to free up rather than shrinking everyone further. The default of 100 means no such sharing is enforced.
Set these per leaf according to who uses it. A service-account leaf such as etl with one submitting user needs a factor high enough to use its own elasticity. A shared exploration leaf needs a minimum percentage so that one person's large query does not starve everyone else.
The application-master trap in small leaves
Each running application needs an application master container before it can ask for workers. maximum-am-resource-percent limits the share of a queue that application masters may occupy; the documented default is 10%, written as 0.1. Its purpose is to stop a flood of submissions from filling a queue with masters that have no room for workers.
In small leaves the arithmetic is unforgiving. Under contention, when bi is held to its 125 GB guarantee, 10% leaves its masters about 12.5 GB. With 2 GB masters, that allows six concurrent applications. The seventh stays in ACCEPTED, and nothing in the job output explains why. Deep trees make this worse, because every level of nesting shrinks leaf capacity while the per-application master size stays the same.
Three fixes, in order of preference: give small interactive leaves a higher percentage, as the example does with 0.3 for bi; merge leaves that are too small to be useful; or reduce the master memory for the engines that run there. See the ResourceManager page for how applications move from ACCEPTED to RUNNING.
ACLs are a union up the tree
This is the rule that catches most clusters. acl_submit_applications and acl_administer_queue are satisfied if the user or group has the permission on the queue or on any of its ancestors. Permissions only add up; a child cannot take away what a parent grants. And the root queue's default ACL is *, which means anyone.
Put those together and a cluster with carefully restricted leaf ACLs, but no root ACL, lets every user submit to every queue. The fix is to close the root by setting its submit ACL to a single space, which means nobody, and then grant access branch by branch, as the configuration above does. ACL values are a comma-separated user list, a space, then a comma-separated group list, which is why " analysts" with a leading space grants a group and no users.
ACLs are only enforced when yarn.acl.enable is true, which is not the default. Test enforcement after every change by submitting as an unauthorised user and confirming the rejection.
Placement: which leaf does a job land in?
Users should not have to know the tree. yarn.scheduler.capacity.queue-mappings maps users or groups to queues with entries of the form u:name:queue or g:name:queue, evaluated in order, first match wins. Placeholders %user, %primary_group and %secondary_group allow patterns such as a per-user leaf under a team parent.
By default a queue named explicitly at submission wins over a mapping. Setting queue-mappings-override.enable to true makes the mapping win instead, which is how you stop users from choosing a better-provisioned queue. Newer releases also accept JSON mapping rules, selected with mapping-rule-format set to json, which add richer matching and an explicit fallback: skip to the next rule, place in the default queue, or reject the application.
Placement interacts with auto-creation. With flexible auto queue creation enabled on a parent, a mapping to a leaf that does not exist can create it from a template, which suits per-user queues. Set the per-parent max-queues limit deliberately; the documented default is 1,000.
Changing the tree without downtime
# Inspect one queue: state, configured and used capacity, running apps
yarn queue -status etl
# Live view of the whole tree as JSON (absolute capacities, usage, user limits)
curl -s "http://$RM:8088/ws/v1/cluster/scheduler" | jq '.scheduler.schedulerInfo.queues'
# Apply an edited capacity-scheduler.xml without restarting the ResourceManager
yarn rmadmin -refreshQueues
# Retire a queue: stop it, wait for it to drain, then remove it and refresh
# 1. set yarn.scheduler.capacity.root.analytics.bi.state = STOPPED, refreshQueues
# 2. wait until the queue shows 0 running and 0 pending applications
# 3. delete bi from root.analytics.queues, rebalance adhoc to 100, refreshQueuesEdit the XML and run yarn rmadmin -refreshQueues. Adding queues and changing capacities, limits and ACLs take effect without restarting the ResourceManager. Removing a queue is stricter: it must first be set to STOPPED, which refuses new submissions while running work finishes, and it must have no running or pending applications. Only then can you delete it and refresh.
Clusters that change queues often can switch to the scheduler configuration mutation API by setting yarn.scheduler.configuration.store.class to leveldb or zk, and change queues through the REST API or the yarn schedulerconf command. Once you do, the XML file is no longer the source of truth, so decide which one is authoritative and stop editing the other.
Keep the configuration in version control either way, and validate a change on a staging ResourceManager first: a refresh that fails validation, for example because siblings no longer sum to 100, is rejected and the old tree stays in place, which is safe but easy to miss.
Design patterns and trade-offs
| Pattern | Shape | Good for | Watch out for |
|---|---|---|---|
| By organisation | root.teamA, root.teamB, leaves per workload | Chargeback, clear ownership | Idle teams hoard guarantees; reorgs force tree changes |
| By service level | root.prod, root.batch, root.adhoc | Protecting SLAs | Weak cost attribution |
| Hybrid | SLA class at level 1, team at level 2 | Most multi-tenant clusters | Depth shrinks leaf size and AM limits |
| Per-user leaves | Auto-created under a team parent | Fair sharing among individuals | Queue count; set max-queues |
Keep the tree shallow, two or three levels below root. Every level multiplies percentages, so deep leaves end up with tiny guarantees, small AM limits and more fragmentation. Guarantees are only as strong as preemption: without it, a queue that lent capacity must wait for borrowed containers to finish. See YARN preemption for how reclaiming works. Hardware partitions such as GPU nodes are better expressed with node labels than with extra queues. If you are choosing between schedulers, the Fair Scheduler page compares the two models and covers migration.
Failure modes
| Symptom | Likely cause | First response |
|---|---|---|
| Jobs stuck in ACCEPTED while the queue is at its guarantee | Leaf's AM limit reached | Raise maximum-am-resource-percent or merge small leaves |
| Single user cannot use idle capacity | user-limit-factor 1 | Raise the factor on that leaf |
| Anyone can submit anywhere | Root ACL left at *, or ACLs disabled | Set root ACL to a space; enable yarn.acl.enable |
| refreshQueues rejected | Sibling percentages do not sum to 100, or mixed modes | Fix the XML; the old tree remains active |
| Queue cannot be removed | Not STOPPED, or apps still pending | Stop, drain, then delete |
| Guaranteed queue waits minutes for capacity | Capacity lent to siblings; preemption off | Enable preemption or cap siblings with maximum-capacity |
What to do next
- Draw your current tree with each leaf's absolute guarantee and maximum, computed as products along the path.
- Check the root submit ACL and
yarn.acl.enable, then test with an unauthorised user. - For every leaf under about 10% of the cluster, compute its AM limit in applications and raise it where interactive users need concurrency.
- Review
user-limit-factorandminimum-user-limit-percentper leaf against who actually submits there. - Move placement into queue mappings with full paths, and decide whether override should be on.
- Put the configuration under version control and practise retiring a queue with STOPPED, drain, delete and refresh.