Troubleshooting starts after someone notices a problem: a query is slow, a dashboard times out. Monitoring is how you notice first. For Apache Impala that means watching three kinds of daemon (impalad coordinators and executors, the statestore, the catalog server) plus the query workload they serve, and turning a few hundred raw metrics into a small set of signals that page a human only when users are, or are about to be, hurt.
This article builds that system. It covers where Impala exposes metrics, how its Prometheus endpoint renames them (with two quirks that break naive queries), which metrics matter for each layer, alert rules you can adapt, query-level service objectives from the query log, and the ways monitoring itself fails. Diagnosing a single bad query, by reading profiles and ExecSummaries, is covered in Impala troubleshooting. Metric names here come from the metric definitions in the current Apache Impala source; older releases may lack some, so confirm against your own endpoint.
The pipeline in one picture
Monitor in layers
A useful monitoring design has layers, each with its own question. Service health: are all daemons up, connected to the statestore, and within memory? Admission and capacity: are queries admitted promptly, or queued, rejected and timed out? Query outcomes: how long do queries take, and how many spill or expire? Metadata: is the catalog propagating changes, or are coordinators running on stale versions? Client front end: can clients connect, or are they waiting for service threads? Alerts belong mainly on admission and outcomes, because those are what users feel. The other layers explain why, and belong on dashboards with a few hard alerts for outright failure.
Know the architecture behind the metrics. The statestore distributes membership and catalog updates to every daemon; the statestore article explains its heartbeats and topics. The catalog server loads metadata and pushes changes; see the catalog service. Admission control decides per pool whether a query may run; see admission control.
Where the metrics live, and how they are renamed
Each daemon serves a debug web UI: by default port 25000 on impalad, 25010 on statestored and 25020 on catalogd. The /metrics page shows metrics for people, and appending ?json returns them as JSON. The page that matters for automation is /metrics_prometheus, which renders metrics in Prometheus text format. It is documented for all three daemons but does not appear in the web UI's list of pages, so people often assume it is missing.
The renaming rules, read from the source, are worth knowing exactly:
- Dots, dashes and quote characters become underscores, and
impala_is prefixed unless already present.impala-server.num-queries-spilledbecomesimpala_server_num_queries_spilled;mem-tracker.process.limitbecomesimpala_mem_tracker_process_limit. - The pool name is part of the metric name, not a label.
admission-controller.total-rejected.root.defaultbecomesimpala_admission_controller_total_rejected_root_default. The dot in the pool name is flattened too, so no selector onroot.defaultwill match; relabel at scrape time to get{pool="root_default"}. - PROPERTY and SET metrics are skipped, so
statestore-subscriber.connectedandstatestore.live-backends.listare not in the Prometheus output. Use the JSON page or other metrics for them. - Time-based units are converted to seconds. Histograms such as
impala-server.query-durations-msare emitted as a summary with_sumand_count, sohistogram_quantile()does not apply. In the current source the quantile labeled 0.2 carries the 25th percentile and the one labeled 0.7 carries the 75th; 0.5, 0.9, 0.95 and 0.999 match their labels. Prefer those, or compute the mean fromrate(_sum) / rate(_count).
Scraping every daemon
Scrape every daemon, label it by role, and turn the pool suffix back into a label so dashboards can group by pool. Prometheus relabelling regexes are anchored, so the list of per-pool families must be explicit; otherwise the first underscore in a pool name gets split in the wrong place.
scrape_configs:
- job_name: impala
metrics_path: /metrics_prometheus
scrape_interval: 30s
static_configs:
- targets: ['coord-1:25000', 'coord-2:25000', 'exec-01:25000', 'exec-02:25000']
labels: {role: impalad}
- targets: ['ss-1:25010']
labels: {role: statestored}
- targets: ['cat-1:25020']
labels: {role: catalogd}
metric_relabel_configs:
# impala_admission_controller_total_rejected_root_default
# -> impala_admission_controller_total_rejected{pool="root_default"}
- source_labels: [__name__]
regex: 'impala_admission_controller_(total_admitted|total_queued|total_rejected|total_timed_out|total_released|total_dequeued|local_num_queued|local_num_admitted_running|local_mem_admitted)_(.+)'
target_label: pool
replacement: '$2'
- source_labels: [__name__]
regex: 'impala_admission_controller_(total_admitted|total_queued|total_rejected|total_timed_out|total_released|total_dequeued|local_num_queued|local_num_admitted_running|local_mem_admitted)_(.+)'
target_label: __name__
replacement: 'impala_admission_controller_$1'If the web UI requires authentication or TLS in your deployment, configure the scrape job to match rather than opening the UI. Keep the target list generated from your inventory, so a new executor cannot run unmonitored.
The signals that matter
| Layer | Impala metric (source name) | What to watch |
|---|---|---|
| Admission | admission-controller.total-rejected.<pool>, total-timed-out.<pool> | any sustained rate is a user-visible failure |
| Admission | admission-controller.local-num-queued.<pool> | queue depth per coordinator; long plateaus mean undersized pools |
| Outcomes | impala-server.query-durations-ms | median and p95 trend per coordinator |
| Outcomes | impala-server.num-queries-spilled | spill rate relative to impala-server.num-queries |
| Outcomes | impala-server.num-queries-expired | idle-timeout expiries; usually clients that never close queries |
| Memory | memory.rss against mem-tracker.process.limit | headroom per impalad |
| Spill disk | tmp-file-mgr.scratch-space-bytes-used | scratch growth toward disk capacity |
| Membership | statestore.live-backends | count against the expected number of daemons |
| Membership | statestore-subscriber.num-connection-failures | increase means a daemon lost the statestore |
| Metadata | catalog.curr-version | should advance on every coordinator after DDL; one that lags its peers is serving stale metadata |
| Metadata | catalog-server.metadata.table.async-loading.queue-len | tables waiting to load |
| Front end | impala.thrift-server.hiveserver2-frontend.connections-in-use | against the configured service thread count |
| Front end | impala.thrift-server.hiveserver2-frontend.connection-setup-queue-size | clients waiting to be accepted |
Alert rules
These rules assume the relabelling above. Thresholds are starting points; set them from two weeks of your own history. Counters reset when a daemon restarts, which increase() and rate() handle.
groups:
- name: impala
rules:
- alert: ImpalaAdmissionRejects
expr: sum by (pool) (increase(impala_admission_controller_total_rejected[10m])) > 0
or sum by (pool) (increase(impala_admission_controller_total_timed_out[10m])) > 0
for: 10m
labels: {severity: page}
annotations: {summary: "Queries rejected or timed out in pool {{ $labels.pool }}"}
- alert: ImpalaQueueBuilding
expr: max by (pool) (impala_admission_controller_local_num_queued) > 10
for: 15m
labels: {severity: ticket}
- alert: ImpalaDaemonMissing
expr: max(impala_statestore_live_backends) < 42 # expected subscribers
for: 5m
labels: {severity: page}
- alert: ImpalaStatestoreFlapping
expr: increase(impala_statestore_subscriber_num_connection_failures[15m]) > 2
labels: {severity: ticket}
- alert: ImpalaMemoryHeadroom
expr: impala_memory_rss / impala_mem_tracker_process_limit > 0.95
for: 10m
labels: {severity: ticket}
- alert: ImpalaSpillRateHigh
expr: sum(rate(impala_server_num_queries_spilled[1h]))
/ sum(rate(impala_server_num_queries[1h])) > 0.2
labels: {severity: ticket}
- alert: ImpalaScrapeDown
expr: up{job="impala"} == 0
for: 5m
labels: {severity: page}Notice what is not paged on: a single slow query, high CPU, or a spill. Those are normal in an analytic engine. Rejections, timeouts, missing daemons and a dead scrape are the failures users and on-call actually need to hear about.
Reading the memory numbers
Memory deserves its own dashboard because Impala's numbers overlap. mem-tracker.process.limit is the limit the daemon enforces on itself; memory.total-used counts TCMalloc and buffer-pool memory; memory.rss is what the operating system sees, including the JVM. A widening gap between RSS and total-used points at JVM or fragmentation growth, which the process limit does not govern and the kernel's out-of-memory killer eventually will. Inside the limit, buffer-pool.reserved against buffer-pool.limit shows how much buffer memory running operators have reserved. Pair the memory panels with tmp-file-mgr.scratch-space-bytes-used per executor: spilling is fine, but scratch disks that fill turn spills into failed queries.
Query-level objectives from the query log
Daemon metrics are aggregated per process; they cannot say "the finance pool's p95 was 40 seconds yesterday". For that, enable workload management, which records every completed query in sys.impala_query_log, and compute service objectives with SQL. Separate queueing time from running time, because they have different owners: queueing belongs to capacity planning, running time to query design.
-- Daily objectives per pool: medians plus the share of queries inside each target.
SELECT to_date(start_time_utc) AS day,
resource_pool,
COUNT(*) AS queries,
APPX_MEDIAN(event_completed_admission) / 1000 AS p50_queue_s,
APPX_MEDIAN(total_time_ms) / 1000 AS p50_total_s,
AVG(CASE WHEN event_completed_admission <= 5000 THEN 1 ELSE 0 END) AS admitted_in_5s,
AVG(CASE WHEN total_time_ms <= 30000 THEN 1 ELSE 0 END) AS done_in_30s
FROM sys.impala_query_log
WHERE start_time_utc >= date_sub(now(), 7)
GROUP BY 1, 2
ORDER BY 1, 2;Here event_completed_admission is the milliseconds from query start until admission finished, so it measures queueing, and total_time_ms the whole query. Expressing the objective as a share, such as 95 percent of BI queries admitted within 5 seconds, avoids needing an exact percentile function and maps directly onto an error budget: if the share falls to 90 percent for a week, the pool needs capacity or tighter limits on its heaviest users. Impala's APPX_MEDIAN gives the median; for other percentiles, export the log to your reporting engine.
A quick check without Prometheus
Before Prometheus is in place, or when it is the thing that broke, a short script can read the endpoint directly. This one prints queue depth and rejections per pool from one coordinator.
import re
import urllib.request
URL = "http://coord-1:25000/metrics_prometheus"
PAT = re.compile(r'^impala_admission_controller_(local_num_queued|total_rejected)_(\S+) (\S+)$')
def pools(url=URL):
out = {}
for line in urllib.request.urlopen(url, timeout=10).read().decode().splitlines():
m = PAT.match(line)
if m:
kind, pool, value = m.groups()
out.setdefault(pool, {})[kind] = float(value)
return out
for pool, v in sorted(pools().items()):
print(f"{pool:30} queued={v.get('local_num_queued', 0):5.0f} rejected_total={v.get('total_rejected', 0):7.0f}")
Worked example: a Monday queue
Monday 08:55, the ImpalaQueueBuilding ticket fires for root_bi: queue depth 25 on both coordinators. No rejections yet. The dashboard shows the pool's admitted memory at its limit, while the other pools sit at half their limits. The executors' RSS is comfortable. That combination means capacity is fine and the pool limit is the bottleneck. The query log shows why: a new dashboard released on Friday runs 14 queries per tile refresh, each reserving 6 GB because its tables have no statistics. The on-call engineer raises the pool's memory limit temporarily, the queue drains in four minutes, and the team computes statistics on the new tables. Per-query reservations fall to under 1 GB, the temporary limit is reverted on Tuesday, and a check is added to the release process. The alert did its job: it fired while users waited seconds, not after they saw timeouts.
Failure modes
- Monitoring only coordinators. Executors hold the memory and the scratch disks; scrape every impalad.
- Querying by pool label without relabelling. Panels show nothing, and people conclude there is no queueing.
- Using
histogram_quantileon summaries, or trusting the 0.2 and 0.7 labels as the 20th and 70th percentiles. - Per-coordinator views read as cluster totals. Queue and admission gauges are local to each coordinator; sum or max them deliberately.
- Alerting on spills and slow queries. On-call learns to ignore pages, and misses the real rejection storm.
- A query log nobody maintains. The Iceberg table grows without snapshot expiry until SLO reports slow down too.
Trade-offs
Scrape interval trades resolution against load: 30 seconds catches queue build-up without hammering the web server. Relabelling at scrape keeps queries readable but must be updated when pools are added, which is a small cost against panels that silently lose a pool. Paging only on user-facing failures reduces fatigue but means capacity problems arrive as tickets, so someone must read them. The query log gives exact per-query objectives at the cost of a table to maintain. Hive users will recognise the same split between JMX metrics and query history in Hive monitoring.
What to do next
- Curl
/metrics_prometheuson one impalad, statestored and catalogd, and confirm the names you plan to use exist in your release. - Add a scrape job that covers every daemon and is generated from your inventory.
- Add the pool relabelling, and verify a pool label appears on admission metrics.
- Load the alert rules, set the expected backend count, and tune thresholds from two weeks of data.
- Build one dashboard per layer: admission per pool, outcomes, memory and spill, membership and metadata, front end.
- Enable workload management and publish a weekly per-pool latency report that separates queueing from running.
- Run a game day: stop an executor and fill a pool, and check that the right alerts fire within minutes.