When an Impala query sits for a minute and then fails with a message about admission timing out, the database did not crash and the query did not run slowly. It never started. It waited in a queue in front of a resource pool, and nothing freed up before its deadline. The query queue is the part of Impala's admission control that users feel most directly, and it is often the least understood: people see queueing as a performance bug when it is the system protecting itself from running out of memory.

This page explains the queue from first principles: why Impala queues instead of just running everything, the three outcomes of an admission decision, how the memory a query needs is worked out, the configuration flags and pool files, how to read queueing in a query profile, and how to tune a pool with a worked example. For the overall design of the admission controller and how it fits into the coordinator and statestore, see Impala admission control architecture.

Why Impala queues queries

Impala is a massively parallel engine that keeps intermediate data such as hash tables, sort buffers and join builds in memory on every executor. A query that runs out of memory on any one host either spills to disk, which slows it down, or fails outright. Now imagine fifty dashboard queries and three large ETL joins arriving in the same second. If all of them start at once, each gets a sliver of memory, many spill, some fail, and the ones that would have finished in two seconds take thirty.

Admission control trades latency at the door for predictability inside. Before a query starts executing, the coordinator checks whether the cluster has room for it under the limits of the resource pool it belongs to. If it does, the query is admitted and runs with the memory it was promised. If not, it waits in that pool's queue until running queries finish and release resources. A short, bounded wait is almost always better than an overcommitted cluster where everything is slow and some things fail.

The queue therefore has three jobs: hold queries that do not fit yet, release them as room appears, and refuse work once waiting no longer makes sense, either because the queue is full or because a query has waited too long.

The admission decision

Clientimpala-shell, JDBCCoordinatorplans queryAdmission checkpool limitsSQLestimateAdmitrun nowRejectqueue fullPool queuewaits, in orderqueueTimed outafter timeoutroom freed: admit headStatestorepool stats topicExecutorsrun fragmentspublish usageadmittedEach coordinator decides locally, using pool usage that every coordinatorpublishes through the statestore, so the picture can be a moment stale.
A query's path through admission: the coordinator plans the query, estimates its memory, and either admits it, queues it in its pool, or rejects it. Queued queries leave by admission or by timing out.

Admission happens on the coordinator after planning and before any fragment starts on an executor. The coordinator compares the query's needs against the limits of its pool and decides one of three outcomes:

  • Admit immediately. There is room now and no earlier query from this coordinator is waiting in the pool, so the query starts.
  • Queue. The query cannot start yet, because the pool is at its running-query limit, its memory limit would be exceeded, a host lacks memory, or other queries are already waiting ahead of it.
  • Reject. The pool's queue is already at its maximum length, the pool is configured to admit nothing, or the query could never fit even on an idle cluster, for example because its per-host memory exceeds what any host can give.

Admission control is decentralized. Every coordinator makes its own decisions, using pool statistics that all coordinators publish through the statestore. The Impala documentation states plainly that this is fast but imprecise under heavy load across many coordinators: the cluster can briefly admit more queries than a limit allows, or hold more in queues than the queue limit, because each coordinator acts on a view that is a heartbeat or so old. The documentation also promises ordering only per submitting host: queries submitted through one coordinator are handled in order, while queries submitted through different hosts have no ordering guarantee. A practical consequence is that a large query at the front of a pool's queue can hold back smaller queries behind it from the same coordinator even when they would fit.

Configuring pools and queues

With no pool files configured, every query goes to a single default pool controlled by impalad startup flags. The defaults below are from the Apache Impala admission configuration reference.

FlagDefaultMeaning
queue_wait_timeout_ms60000How long a request waits to be admitted before it times out
default_pool_max_requests-1 (unlimited)Running queries allowed before new ones queue
default_pool_max_queuedunlimitedQueued queries allowed before new ones are rejected
default_pool_mem_limitempty (unlimited)Cluster-wide memory all running queries in the pool may use
fair_scheduler_allocation_pathemptyPath to fair-scheduler.xml defining pools
llama_site_pathemptyPath to llama-site.xml with per-pool limits

Real clusters define several pools. fair-scheduler.xml declares the pools, their maximum memory and who may submit to them; llama-site.xml holds the queue-related limits per pool. A trimmed example for a dashboard pool:

<!-- fair-scheduler.xml -->
<allocations>
  <queue name="root">
    <queue name="dashboards">
      <maxResources>200000 mb, 0 vcores</maxResources>
      <aclSubmitApps>bi_users</aclSubmitApps>
    </queue>
    <queue name="etl">
      <maxResources>600000 mb, 0 vcores</maxResources>
      <aclSubmitApps>etl_svc</aclSubmitApps>
    </queue>
  </queue>
  <queuePlacementPolicy>
    <rule name="specified" create="false"/>
    <rule name="default"/>
  </queuePlacementPolicy>
</allocations>

<!-- llama-site.xml -->
<configuration>
  <property>
    <name>llama.am.throttling.maximum.placed.reservations.root.dashboards</name>
    <value>20</value>   <!-- max running -->
  </property>
  <property>
    <name>llama.am.throttling.maximum.queued.reservations.root.dashboards</name>
    <value>50</value>   <!-- max queued -->
  </property>
  <property>
    <name>impala.admission-control.pool-queue-timeout-ms.root.dashboards</name>
    <value>30000</value>
  </property>
</configuration>

Clients choose a pool with the REQUEST_POOL query option, for example SET REQUEST_POOL=root.dashboards; in impala-shell, or through your driver's query-option settings. The placement policy above honours a specified pool if the user may submit to it and otherwise falls back to the default pool.

How memory to admit is computed

The memory side of admission is where most queueing surprises come from. For each query the coordinator settles on a per-host memory amount to admit against. If the user or the pool defaults set MEM_LIMIT, that value is used. Otherwise the planner's per-host estimate is used, which depends heavily on table and column statistics. Pools can bound the result with max-query-mem-limit and min-query-mem-limit, and clamp-mem-limit-query-option controls whether a user-supplied MEM_LIMIT is also clamped into that range.

The coordinator then checks two things. At the pool level, the per-host amount multiplied by the number of hosts the query will run on, added to the memory already admitted in the pool, must stay under the pool's maximum memory. At the host level, each executor the query will use must have enough admittable memory left. Failing either check queues the query.

This explains a classic pattern. A table without statistics produces a wildly inflated estimate, so a query that needs 500 MB per host is admitted against 20 GB per host, and three such queries fill a pool that should hold forty. The fix is not a bigger pool but COMPUTE STATS on the tables involved, described in Impala statistics and metadata, plus sensible minimum and maximum per-query limits. The opposite mistake, underestimating, lets queries in that later spill or run out of memory, which is covered in Impala memory limits.

Worked example: sizing a dashboard pool

Take a cluster of ten executors, each with about 100 GB admittable to queries. The dashboards pool above allows 200 GB across the cluster, 20 running queries, 50 queued and a 30-second timeout, with per-query limits between 1 GB and 8 GB per host.

A typical dashboard query, with good statistics, is estimated at 1.5 GB per host and runs on all ten executors, so it is admitted against 15 GB of the pool. The pool's memory allows 200 / 15, or 13 such queries at once, which is below the 20-query limit, so memory is the binding constraint. If each query takes four seconds, the pool completes about 13 / 4, or roughly 3.25 queries per second at steady state. During a burst where 60 dashboards refresh together, 13 start, 47 queue, and the last of them waits about 47 / 3.25, or roughly 15 seconds, which is inside the 30-second timeout. Nothing is rejected.

Now someone drops statistics on a large dimension table. The estimate jumps, gets clamped to the 8 GB per-host maximum, and each query is admitted against 80 GB. Only two run at once, throughput falls to about half a query per second, and in the same burst most queries hit the timeout. The profile and the arithmetic point at the same culprit, and the query plan in EXPLAIN shows the inflated per-host estimate.

-- In impala-shell, for one slow dashboard query:
SET REQUEST_POOL=root.dashboards;
EXPLAIN SELECT ...;            -- check the per-host memory estimate
SELECT ...;
PROFILE;                       -- then search the output for "Admission result"

Reading queueing in profiles and metrics

Every query profile records how admission went. The key line is Admission result, whose value is one of Admitted immediately, Admitted (queued), Rejected or Timed out (queued). Next to it the profile records the reason a query was queued or rejected, in wording taken from the admission controller, such as the number of running queries being over the pool's limit, not enough memory in the pool or on a host, or queue full, limit=50, num_queued=50 for a rejection. The query timeline shows when the query was submitted, when it was queued and when admission completed, so the gap between those events is the time spent waiting.

Beyond single queries, the impalad web UI and its metrics show per-pool counts of running, queued, rejected and timed-out queries, plus memory admitted versus the limit. Track these over time. A pool that queues briefly at peaks is working as intended; a pool whose queue rarely empties, or whose timeout count climbs, needs attention.

Failure modes

  • Timeouts at peak. Bursts exceed what the pool can drain within the timeout. Raise throughput by fixing estimates or adding capacity, or spread refresh schedules, rather than simply lengthening the timeout.
  • Head-of-line blocking. One large query at the front of the queue waits for lots of memory while small queries behind it, which would fit, wait too. Put large and small workloads in separate pools.
  • Inflated estimates. Missing or stale statistics make queries reserve far more memory than they use, so pools fill with few queries. Compute statistics and set per-query maximums.
  • Overadmission across coordinators. Many coordinators acting on slightly stale shared state can briefly admit beyond a limit. Leave headroom, or route a pool's traffic through fewer coordinators.
  • Retry storms. Clients that resubmit immediately on rejection keep the queue full. Retry with backoff and jitter, and surface rejections to users instead of hiding them.
  • Wrong pool. Queries land in the default pool because REQUEST_POOL was not set or the user lacks submit rights. Check the pool name in the profile.

Tuning and trade-offs

KnobRaise it whenRisk if too high
Pool max memoryQueries queue on memory while hosts sit idlePools starve one another; more spilling
Max runningMemory is free but the count limit bindsCPU contention slows every query
Max queuedShort bursts are rejected outrightUsers wait longer only to time out anyway
Queue timeoutBatch jobs can tolerate waitingInteractive users stare at a spinner
Per-query max memoryLegitimate large queries spill or failA few queries monopolise the pool

The honest trade-off is between utilization and predictability. Tight limits keep latency consistent and protect against memory failures but leave capacity unused off-peak; loose limits raise utilization but bring back the overcommitted cluster that admission control exists to prevent. Separate pools for interactive and batch work, each sized to its own pattern, usually beat one large tuned pool.

What to do next

  1. List your pools and record max memory, max running, max queued and queue timeout for each, from the pool files or the impalad flags.
  2. Pull a week of profiles or metrics and count queries that were admitted after queueing, rejected or timed out, per pool.
  3. For the slowest queued queries, compare the per-host estimate in EXPLAIN with actual peak memory, and run COMPUTE STATS where they differ widely.
  4. Separate interactive and batch work into different pools if they share one now.
  5. Make clients set REQUEST_POOL explicitly and retry rejections with backoff.
  6. Read Impala troubleshooting for the wider diagnostic workflow.
Key takeaway: The Impala query queue is the price of a cluster that does not run out of memory. Size each pool from the per-host memory its queries really need, keep statistics fresh so estimates stay honest, separate interactive and batch work, and read the Admission result line in the profile before you change any limit.