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

Impala monitoring pipelineimpaladweb UI :25000statestoredweb UI :25010catalogdweb UI :25020Prometheusscrape, relabel, rules/metrics_prometheusAlertmanagerpage or ticketDashboardsper layer, per poolsys.impala_query_logper-query historyworkload managementSLO reportslatency per pool
Each daemon serves Prometheus text on its web UI port; Prometheus scrapes, relabels pools and evaluates rules for Alertmanager and dashboards, while workload management records each query in an Iceberg table for SLO reports.

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-spilled becomes impala_server_num_queries_spilled; mem-tracker.process.limit becomes impala_mem_tracker_process_limit.
  • The pool name is part of the metric name, not a label. admission-controller.total-rejected.root.default becomes impala_admission_controller_total_rejected_root_default. The dot in the pool name is flattened too, so no selector on root.default will match; relabel at scrape time to get {pool="root_default"}.
  • PROPERTY and SET metrics are skipped, so statestore-subscriber.connected and statestore.live-backends.list are 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-ms are emitted as a summary with _sum and _count, so histogram_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 from rate(_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

LayerImpala metric (source name)What to watch
Admissionadmission-controller.total-rejected.<pool>, total-timed-out.<pool>any sustained rate is a user-visible failure
Admissionadmission-controller.local-num-queued.<pool>queue depth per coordinator; long plateaus mean undersized pools
Outcomesimpala-server.query-durations-msmedian and p95 trend per coordinator
Outcomesimpala-server.num-queries-spilledspill rate relative to impala-server.num-queries
Outcomesimpala-server.num-queries-expiredidle-timeout expiries; usually clients that never close queries
Memorymemory.rss against mem-tracker.process.limitheadroom per impalad
Spill disktmp-file-mgr.scratch-space-bytes-usedscratch growth toward disk capacity
Membershipstatestore.live-backendscount against the expected number of daemons
Membershipstatestore-subscriber.num-connection-failuresincrease means a daemon lost the statestore
Metadatacatalog.curr-versionshould advance on every coordinator after DDL; one that lags its peers is serving stale metadata
Metadatacatalog-server.metadata.table.async-loading.queue-lentables waiting to load
Front endimpala.thrift-server.hiveserver2-frontend.connections-in-useagainst the configured service thread count
Front endimpala.thrift-server.hiveserver2-frontend.connection-setup-queue-sizeclients 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_quantile on 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

  1. Curl /metrics_prometheus on one impalad, statestored and catalogd, and confirm the names you plan to use exist in your release.
  2. Add a scrape job that covers every daemon and is generated from your inventory.
  3. Add the pool relabelling, and verify a pool label appears on admission metrics.
  4. Load the alert rules, set the expected backend count, and tune thresholds from two weeks of data.
  5. Build one dashboard per layer: admission per pool, outcomes, memory and spill, membership and metadata, front end.
  6. Enable workload management and publish a weekly per-pool latency report that separates queueing from running.
  7. Run a game day: stop an executor and fill a pool, and check that the right alerts fire within minutes.
Key takeaway: Monitor Impala in layers and page only on what users feel. Scrape /metrics_prometheus on every impalad, statestored and catalogd; remember that names are flattened with underscores and an impala_ prefix, that pool names are embedded in admission metric names until you relabel them, that property and set metrics are absent, and that durations are summaries rather than histograms. Alert on rejections, timeouts, missing daemons and dead scrapes; put memory, spill, membership, metadata and front-end signals on dashboards; and use the workload-management query log for per-pool objectives that separate queueing from running.