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

1. Startup flagsflag file per role; restart2. Pool definitionsfair-scheduler.xml3. Pool settingsllama-site.xml4. Session optionsSET in client or JDBC URLCoordinatoradmission + planningprocess defaultspool, limitspool defaults, clampsSET overridesEffective optionsin the query profileFour configuration layers meet in the coordinator; the profile records the result for every query.
Flags, two pool files and session options all feed the coordinator; the query profile shows what won.

Who owns each layer

LayerWhere it livesWho changes itTakes effect
Startup flagsCommand line or a flag file per daemon role, usually rendered by a management toolPlatform teamDaemon restart
Pool definitionsfair-scheduler.xml, path set by --fair_scheduler_allocation_pathPlatform team with tenant ownersRead by the impalads; confirm on your version how quickly edits are picked up
Pool settingsllama-site.xml, path set by --llama_site_pathPlatform teamSame as pool definitions
Session optionsSET statements, client config, connection URLUsers and applicationsNext 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.bi with 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 maxResources of 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.

  1. Change in one place. Templates for flag files and both pool files live in a repository; a change is a reviewed commit.
  2. Validate before deploy. Parse both XML files, check every llama-site.xml property ends in a pool that exists in fair-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.
  3. 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.
  4. Verify the running state. Read /varz on every daemon and the /admission page on a coordinator, and diff them against the intended values.
  5. 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.xml property for root.Bi does not configure root.bi. Validate names against the pool file.
  • Byte units. Memory bounds are bytes; writing 4g or forgetting three zeros gives a pool that admits nothing or everything.
  • Clamp off by accident. Without the clamp, one SET MEM_LIMIT in 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

  1. Put flag files and both pool files in version control, one flag template per role.
  2. Pull /varz?json from every daemon and fix any per-host differences.
  3. Make root.default small and route known service accounts to pools by group, not by client settings.
  4. Set min and max query memory limits with the clamp on for shared pools, sized from real profile peaks.
  5. Check that the sum of pool budgets versus total executor memory is a deliberate choice.
  6. Add the validation checks above to the deploy pipeline and run the drift script nightly.
  7. Teach users to read the Query Options (set by configuration) line in a profile before filing a ticket.
Key takeaway: Impala behaviour comes from four layers: startup flags, the pool definition file, per-pool settings and session options. Keep everything that protects the cluster in layers users cannot change, clamp pool memory bounds, size pool budgets with explicit arithmetic, keep every file in version control, verify the running state through /varz and /admission, and read the profile's query options line to see what a query actually used.