Ask why an Impala query ran with an 8 GB memory limit, timed out after 30 minutes, or waited in a queue, and the answer usually lives in one of four places: a startup flag on the coordinator, a resource pool file, a pool's default query options, or a SET statement in the client session. Each layer has a different owner, a different change process and a different way to verify it. Most configuration incidents are not a wrong value; they are a right value in the wrong layer, or two layers disagreeing.
This article maps those layers, shows real configuration for each, explains how to confirm which value a query actually used, and gives a safe change process with drift detection. The daemon roles, the important individual flags, ports and rolling restarts are covered in Impala administration and the admission algorithm itself in admission control. This page is about the configuration system around them.
The four layers
Who owns each layer
| Layer | Where it lives | Who changes it | Takes effect |
|---|---|---|---|
| Startup flags | Command line or a flag file per daemon role, usually rendered by a management tool | Platform team | Daemon restart |
| Pool definitions | fair-scheduler.xml, path set by --fair_scheduler_allocation_path | Platform team with tenant owners | Read by the impalads; confirm on your version how quickly edits are picked up |
| Pool settings | llama-site.xml, path set by --llama_site_path | Platform team | Same as pool definitions |
| Session options | SET statements, client config, connection URL | Users and applications | Next query in that session |
Two rules follow from this table. First, anything that protects the cluster (memory caps, timeouts, concurrency) must live in a layer users cannot change, which means layers 1 to 3. Second, every layer must be in version control, because each one changes behaviour and only the session layer leaves a trace in the query that used it.
Startup flags and flag files
Impala daemons read command-line flags, and a flag file is simply a list of the same flags, one per line, passed with --flagfile. Management tools such as Cloudera Manager render the file from their own settings and overwrite hand edits on the next deploy, so change flags in the tool, or in the template your automation renders, never on the host.
Keep one flag file per role, not per host. Coordinators, executors, catalogd and statestored need different flags, and per-host differences are almost always drift rather than design. Hardware differences should be expressed as a separate role (for example executors with large memory) rather than as edits to individual hosts.
# coordinator.flags (rendered from a template; do not edit on the host)
--is_executor=false
--mem_limit=48g
--default_query_options=query_timeout_s=1800,mem_limit=4g
--idle_session_timeout=3600
--idle_query_timeout=600
--fair_scheduler_allocation_path=/etc/impala/conf/fair-scheduler.xml
--llama_site_path=/etc/impala/conf/llama-site.xml--default_query_options takes comma-separated key=value pairs and sets the process-wide starting value for query options. It is the bluntest layer: it applies to every pool on that coordinator, and it is invisible to users until they read a profile. Use it for genuinely universal defaults such as a query timeout, and put anything that differs by workload into pools instead. The effective flag values of a running daemon are on its web UI at /varz; that page, not your template, is the truth.
Resource pools: the two files
Resource pools are defined in a file that borrows YARN's Fair Scheduler format. Impala uses it for pool names, hierarchy, the aggregate memory a pool may use across the cluster, submission ACLs and the rules that place a query in a pool. CPU figures in the file are not what Impala admits on; memory is.
<allocations>
<queue name="root">
<queue name="bi">
<maxResources>800000 mb, 0 vcores</maxResources>
<aclSubmitApps> bi_users</aclSubmitApps>
</queue>
<queue name="etl">
<maxResources>1200000 mb, 0 vcores</maxResources>
<aclSubmitApps>etl_svc etl_admins</aclSubmitApps>
</queue>
<queue name="default">
<maxResources>200000 mb, 0 vcores</maxResources>
</queue>
</queue>
<queuePlacementPolicy>
<rule name="specified" create="false"/>
<rule name="default"/>
</queuePlacementPolicy>
</allocations>The ACL value follows Hadoop's format: users, then a space, then groups, so bi_users with a leading space means 'no users, the group bi_users'. The placement policy shown honours an explicit REQUEST_POOL option (specified) and otherwise falls through to root.default. Keep the default pool small, so an unconfigured client cannot consume the cluster.
Per-pool settings live in the second file, as properties whose names end with the pool name. These are the ones documented for admission control:
<configuration>
<!-- concurrency: running and queued queries for root.bi -->
<property><name>llama.am.throttling.maximum.placed.reservations.root.bi</name><value>20</value></property>
<property><name>llama.am.throttling.maximum.queued.reservations.root.bi</name><value>100</value></property>
<property><name>impala.admission-control.pool-queue-timeout-ms.root.bi</name><value>120000</value></property>
<!-- per-host memory bounds and whether they override a user's MEM_LIMIT -->
<property><name>impala.admission-control.min-query-mem-limit.root.bi</name><value>1073741824</value></property>
<property><name>impala.admission-control.max-query-mem-limit.root.bi</name><value>4294967296</value></property>
<property><name>impala.admission-control.clamp-mem-limit-query-option.root.bi</name><value>true</value></property>
<!-- defaults applied to every query in the pool -->
<property><name>impala.admission-control.pool-default-query-options.root.bi</name>
<value>query_timeout_s=600,num_scanner_threads=8</value></property>
</configuration>The memory values are bytes. The placed.reservations and queued.reservations names are historical (they come from Llama, a long-retired component), but they are still the documented keys for maximum running and maximum queued queries.
Query options and memory clamps
A query's options are assembled from the process defaults, the pool's default options and whatever the session set. Users override pool defaults with SET; that is by design, because analysts legitimately tune options per query. The memory bounds are the exception that makes pools safe. The Impala documentation describes them this way: when min and max query memory limits are set, admission control chooses a memory limit between them based on the query's per-host estimate, and if the clamp setting is true and the user sets MEM_LIMIT outside that range, the effective limit becomes the minimum or the maximum.
So with the pool above, a user who runs SET MEM_LIMIT=64g in root.bi still gets 4 GB per host. Set the clamp to false and the user's value wins, which is occasionally right for a trusted ETL pool and almost never right for a BI pool. Other cluster-protecting options in pool defaults, such as a query timeout, can still be overridden by a session, so do not rely on them alone; pair them with the daemon flags for idle sessions and queries.
The documentation does not spell out every interaction between --default_query_options and pool defaults, so do not reason about it; measure it. Every query profile contains a Query Options (set by configuration) line listing the non-default options the query ran with, and a Query Options (set by configuration and planner) line after planning. Those lines are the final answer for that query.
-- In impala-shell: which pool and options will my next query use?
SET REQUEST_POOL=root.bi;
SET ALL; -- every option and its current session value
SELECT count(*) FROM sales.orders WHERE order_date = '2026-10-01';
PROFILE; -- then search for "Query Options (set by configuration"
Pool memory arithmetic
Admission control budgets memory per pool across the cluster, so the numbers in the two files must be designed together. Take a cluster of 20 executors, each started with an 80 GB process memory limit, so 1,600 GB in total.
- A query admitted in
root.biwith the 4 GB per-host maximum and running on all 20 executors is accounted as up to 4 GB x 20 = 80 GB of the pool's budget. - The pool's
maxResourcesof 800,000 MB (about 781 GB) therefore admits about 9 such queries at their maximum before queueing, even though the running-query limit says 20. Small queries with lower estimates fit more. - The ETL pool's 1,200,000 MB plus BI's 800,000 MB plus default's 200,000 MB add up to more than the cluster's 1,600 GB. That oversubscription is deliberate: pools rarely peak together, and each host's own limit is still enforced. If they do peak together, queries queue on host memory instead of on pool budgets, which is harder to explain to users. Decide which behaviour you want and write it down.
Get the per-host maximum from real profiles, not guesses: take the peak per-host memory of the pool's queries over a few weeks (memory limits explains where it appears), set the maximum around the 95th to 99th percentile, and let the outliers spill or fail visibly.
Changing configuration safely
Treat configuration as code with a deploy pipeline, even if the deploy step is a management-tool API call.
- Change in one place. Templates for flag files and both pool files live in a repository; a change is a reviewed commit.
- Validate before deploy. Parse both XML files, check every
llama-site.xmlproperty ends in a pool that exists infair-scheduler.xml(do not count on a typo in a pool name producing a visible error), and check memory bounds are below the executors' process limit. - Canary. Restart one coordinator first for flag changes, route a small share of traffic to it, and compare its error and queue rates with the others.
- Verify the running state. Read
/varzon every daemon and the/admissionpage on a coordinator, and diff them against the intended values. - Record. Note the commit in a change log the on-call engineer can see, because a configuration change is the first suspect in any next-day slowdown.
The drift check is a short script against the web UI, which accepts ?json on every page:
import json, sys, urllib.request
WANTED = {"mem_limit": "48g", "idle_session_timeout": "3600",
"default_query_options": "query_timeout_s=1800,mem_limit=4g"}
def walk(node):
"""Yield (name, value) pairs from whatever JSON shape this version uses."""
if isinstance(node, dict):
if "name" in node and "value" in node:
yield str(node["name"]), str(node["value"])
for v in node.values():
yield from walk(v)
elif isinstance(node, list):
for v in node:
yield from walk(v)
drift = 0
for host in sys.argv[1:]: # coordinator hosts only
url = f"http://{host}:25000/varz?json" # add TLS and auth if enabled
flags = dict(walk(json.load(urllib.request.urlopen(url, timeout=10))))
for k, want in WANTED.items():
if flags.get(k) != want:
drift += 1
print(f"{host}: {k}={flags.get(k)!r}, expected {want!r}")
sys.exit(1 if drift else 0)The walker does not assume a particular JSON layout, because the page's structure has varied between releases; check its output against /varz in a browser once on your version. If the web UI is protected, as it should be (see the shell and web UI guide), pass credentials and use HTTPS. Run it from the deploy pipeline and nightly, and alert on a non-zero exit.
Worked example: the dashboards in the wrong pool
Worked example. After a release, dashboards in the BI pool start failing with memory-limit errors at 2 GB per host, although the pool's maximum is 4 GB. The on-call engineer opens a failed query's profile. The set by configuration line shows MEM_LIMIT=2147483648 and REQUEST_POOL=root.default, not root.bi.
So the queries are not in the BI pool at all. The BI tool's connection URL used to include the pool option; the new release of the tool rebuilt the connection string and dropped it, and the placement policy fell through to root.default, whose pool defaults set a 2 GB limit. Nothing in Impala changed. The fix is to restore the pool option in the client and to add a placement rule based on the BI service account's group (the Fair Scheduler format supports rules such as primaryGroup), so routing no longer depends on a client setting. Test any new rule against your Impala version before rollout. The lesson: put routing for shared service accounts in the pool file, where a client upgrade cannot remove it.
Failure modes
- Typo in a pool name. A
llama-site.xmlproperty forroot.Bidoes not configureroot.bi. Validate names against the pool file. - Byte units. Memory bounds are bytes; writing
4gor forgetting three zeros gives a pool that admits nothing or everything. - Clamp off by accident. Without the clamp, one
SET MEM_LIMITin a notebook bypasses the pool's maximum. - Hand edits on hosts. The management tool overwrites them on the next deploy, reverting a fix at the worst moment.
- Per-host drift. One coordinator with an old flag file plans differently; users see 'random' failures depending on which coordinator the load balancer chose.
- Over-tight queue timeouts. Queries fail with queue timeouts while executors are idle, because the pool budget, not the hardware, is the bottleneck.
Trade-offs
More pools give better isolation and clearer chargeback, but fragment memory: each pool's idle budget cannot help a busy neighbour beyond the oversubscription you allow. Clamping protects the cluster but takes control away from experienced users, which pushes them to ask for their own pools. Process-wide defaults are simple but invisible; pool defaults are visible in the pool file but easy to bypass with SET. A reasonable split: safety limits in flags and clamped pool bounds, workload tuning in pool defaults, and per-query tuning left to users. When something still goes wrong, Impala troubleshooting has the symptom-by-symptom playbook.
What to do next
- Put flag files and both pool files in version control, one flag template per role.
- Pull
/varz?jsonfrom every daemon and fix any per-host differences. - Make
root.defaultsmall and route known service accounts to pools by group, not by client settings. - Set min and max query memory limits with the clamp on for shared pools, sized from real profile peaks.
- Check that the sum of pool budgets versus total executor memory is a deliberate choice.
- Add the validation checks above to the deploy pipeline and run the drift script nightly.
- Teach users to read the Query Options (set by configuration) line in a profile before filing a ticket.