Impala is a massively parallel SQL engine that reads tables registered in the Hive Metastore directly from HDFS, S3 or Ozone, without MapReduce or Tez in between. Its speed comes from long-running daemons that keep metadata cached and execute query fragments in memory. That design is also what makes it an administration job: there are three kinds of daemon, each with its own flags, ports, logs and failure modes, and the metadata cache has to be kept in step with a metastore that other engines change behind Impala's back.

This article is a runbook for that job. It covers what each daemon does, the flags that matter, the ports and web pages you will use every day, high availability for the statestore and catalog, draining nodes for rolling restarts, metadata and statistics routines, security wiring and monitoring. It deliberately stays at the operational level. The internals of admission control, memory limits and the catalog have their own pages, and so does query troubleshooting.

Advertisement

The daemons and what you are responsible for

impalad is the worker and, optionally, the coordinator. A coordinator accepts client connections, parses and plans queries, runs admission control and assembles results. An executor runs plan fragments: scans, joins, aggregations and sorts. The same binary does both by default, and on anything beyond a small cluster you should split them with --is_executor=false on a few coordinators and --is_coordinator=false on the rest. Coordinators then carry the session and planning load, and executors can be sized, restarted and scaled without dropping client connections.

statestored tracks cluster membership and broadcasts topics, including the catalog topic and admission control state. If it is down, running queries continue, but membership changes and metadata updates stop propagating. catalogd loads table metadata from the Hive Metastore and file listings from storage, and pushes changes out through the statestore. If it is down, queries against already-cached tables keep working, but DDL and metadata refreshes fail.

What an Impala administrator runs, and what each piece depends onClientsimpala-shell, JDBC/ODBC, BICoordinatorsimpalad is_executor=falseExecutorsimpalad is_coordinator=falsestatestoredmembership, topics (HA pair)catalogdmetadata cache (HA pair)Hive Metastoreschemas, partitionsHDFS / S3 / Ozonedata filesRanger, Kerberosauthz and authnHS2 21050KRPC 27000subscribecatalog topicHMS APIscanpoliciesWeb UIs: impalad 25000, statestored 25010, catalogd 25020. Every box needs its own flags, logs and monitoring.
Clients talk to coordinators, coordinators fan out to executors, and both subscribe to the statestore. catalogd sits between the Hive Metastore and the cluster's metadata caches.

The flags that matter

Impala daemons are configured with command-line flags, usually through a management tool that renders them into a flag file. A handful of flags decide most of a cluster's behaviour:

# Executor (the many): runs fragments, owns the memory and scratch disks
--is_coordinator=false
--mem_limit=80%                                   # of physical RAM; leave room for the OS and other services
--scratch_dirs=/data/1/impala/scratch:500GB,/data/2/impala/scratch:500GB
--state_store_host=ss1.example.com
--catalog_service_host=cat1.example.com
--use_local_catalog=true                          # fetch metadata on demand instead of the full topic

# Coordinator (the few): plans queries, holds sessions, runs admission control
--is_executor=false
--mem_limit=64g
--default_query_options=mem_limit=8g,query_timeout_s=1800
--idle_session_timeout=3600
--use_local_catalog=true

# catalogd
--catalog_topic_mode=minimal                      # pairs with use_local_catalog on the impalads
--hms_event_polling_interval_s=2                  # pick up Hive/Spark DDL without manual INVALIDATE

The memory limit is the single most important number. It caps what Impala will allocate on a host, and admission control budgets queries against it. Setting it as a percentage is convenient on uniform hardware; on shared nodes set an absolute value that leaves room for the DataNode, the OS page cache and anything else co-located. Scratch directories are where operators spill when a query exceeds its memory; put them on local SSDs, one per disk, and cap each so a runaway spill cannot fill a disk that also holds data.

The local catalog pair, --use_local_catalog=true on impalads and --catalog_topic_mode=minimal on catalogd, makes coordinators fetch metadata on demand and cache it with eviction, instead of each holding a full copy of every table. On clusters with thousands of tables or partitions, it is the difference between coordinators with small heaps and coordinators that need very large ones. Setting --hms_event_polling_interval_s to a positive value makes catalogd follow the metastore's notification log, so tables created or altered by Hive or Spark appear without manual invalidation.

Advertisement

Ports and web pages

PortDaemonUsed for
21050impalad (coordinator)HiveServer2 protocol: JDBC, ODBC, impala-shell
21000impalad (coordinator)Legacy Beeswax protocol
28000impalad (coordinator)HiveServer2 over HTTP, for load balancers and proxies
27000impaladKRPC between daemons (Impala 3.2 and later)
25000impaladDebug web UI
24000 / 25010statestoredSubscriber service / web UI
26000 / 25020catalogdCatalog service / web UI

The web UIs are the fastest diagnostic tool you have, and every page accepts ?json for scripting. On an impalad, /queries lists running and recent queries with links to their profiles, /sessions shows open sessions, /backends shows cluster membership as this coordinator sees it, /memz breaks down memory by tracker, /admission shows pools and queues, /varz shows effective flag values and /metrics exposes counters. On statestored, /subscribers and /topics show who is connected and how large each topic is. On catalogd, /catalog lists the objects it holds. Restrict these ports to admin networks, or enable authentication on them, because they reveal query text and configuration.

High availability for statestore and catalog

Current Impala releases can run the statestore and catalogd as active-standby pairs. The two features are independent: you can enable either or both. With catalog HA, a standby catalogd is promoted when the active one fails. With statestore HA, two statestores watch each other and subscribers are configured with both.

# catalogd HA: both catalogds AND the statestore(s)
--enable_catalogd_ha=true

# statestored HA: on each of the two statestores
--enable_statestored_ha=true
--state_store_ha_port=24020                       # example port; pick one free on both hosts
--state_store_peer_host=ss2.example.com            # the other statestore
--state_store_peer_ha_port=24020
# then restart every subscriber (catalogd, coordinators, executors) with
# enable_statestored_ha=true and flags naming BOTH statestores; check the
# HA page for your release for the exact subscriber flag names.

Two operational notes. First, a promoted catalogd starts with a cold cache, so the first queries against each table after failover load metadata and are slower; this is still far better than an outage. Second, HA does not protect against a Hive Metastore outage. Put the metastore and its database behind their own HA, as described in the Hive Metastore guides, or Impala will keep working only for already-cached metadata. The coordinator tier gets its availability differently: run at least two coordinators behind a load balancer that supports session affinity, because an HS2 session lives on one coordinator.

Graceful shutdown and rolling restarts

Restarting an executor mid-query fails every query with a fragment on it. The :shutdown() command avoids that. Issued through impala-shell, it tells the target daemon to announce that it is shutting down, wait a grace period so coordinators stop scheduling new work on it, and exit once its running fragments complete or --shutdown_deadline_s expires; the deadline defaults to one hour. Executors can be drained this way without disrupting running queries. Coordinators are different: queries submitted to a coordinator after its shutdown starts fail, so take a coordinator out of the load balancer first and then shut it down.

#!/usr/bin/env bash
# Rolling restart of executors, one at a time, draining each with :shutdown().
set -euo pipefail
COORD=coord1.example.com
for host in $(cat executors.txt); do
  echo "draining $host"
  impala-shell -i "$COORD" -q ":shutdown('${host}:27000')"      # KRPC port on Impala 3.2+
  # wait until the backend disappears from the coordinator's membership view
  until ! curl -s "http://${COORD}:25000/backends?json" | grep -q "\"${host}:"; do sleep 10; done
  ssh "$host" 'sudo systemctl restart impalad'
  until curl -s "http://${COORD}:25000/backends?json" | grep -q "\"${host}:"; do sleep 10; done
  echo "$host back in membership"; sleep 60                     # let caches warm before the next one
done

The script restarts executors one at a time and waits for each to rejoin membership before moving on. With a data cache enabled, give each node a little time to warm before draining the next, or scan latency degrades across the whole restart. For upgrades, the order that avoids protocol mismatches is the one your distribution documents; the usual pattern is statestore and catalog first, then executors, then coordinators, within one maintenance window.

Metadata hygiene

Most confusing Impala behaviour is stale metadata: a table that other engines changed, and whose files or schema Impala has not seen. Event processing handles DDL and most inserts from Hive and Spark. When it is off, or for writers that bypass the metastore, you need two statements and should know which one to use:

-- New files appended to an existing table or partition by Spark or an ingest job
REFRESH sales.orders;
REFRESH sales.orders PARTITION (dt='2026-09-30');      -- much cheaper on large tables

-- A table created, dropped or restructured outside Impala, with event processing off
INVALIDATE METADATA sales.orders;                      -- always name the table

-- Statistics after a large load; incremental stats only for partitions that changed
COMPUTE STATS sales.customers;
COMPUTE INCREMENTAL STATS sales.orders PARTITION (dt='2026-09-30');
SHOW TABLE STATS sales.orders;                          -- #Rows = -1 means no stats

REFRESH reloads file listings for an existing table and is cheap, especially when limited to a partition. INVALIDATE METADATA discards everything catalogd knows about a table and forces a full reload on next access; it is the right tool for tables created or changed structurally outside Impala. Never run it without a table name on a production cluster: that discards the whole catalog and every table reloads on its next use.

Statistics drive the planner's join order and join strategy. A table with no stats shows -1 rows in SHOW TABLE STATS, and a query joining two such tables can choose a broadcast join that ships a huge table to every node. Schedule COMPUTE INCREMENTAL STATS after each partition load, and full COMPUTE STATS for unpartitioned tables after large changes.

Security wiring

A production cluster needs authentication, encryption and authorization, and each has its own flags:

# Kerberos (all daemons)
--principal=impala/_HOST@EXAMPLE.COM
--keytab_file=/etc/impala/impala.keytab

# TLS for client and internal connections
--ssl_server_certificate=/etc/impala/tls/server.pem
--ssl_private_key=/etc/impala/tls/server.key
--ssl_client_ca_certificate=/etc/impala/tls/ca.pem

# LDAP for BI users who do not have Kerberos tickets (coordinators)
--enable_ldap_auth=true
--ldap_uri=ldaps://ldap.example.com

# Ranger authorization (impalads and catalogd)
--server_name=server1
--authorization_provider=ranger
--ranger_service_type=hive
--ranger_app_id=impala

Kerberos authenticates both clients and daemons to each other. LDAP is typically enabled only on coordinators, for BI tools and users who do not hold Kerberos tickets, and must be combined with TLS because passwords travel in the connection. Ranger policies are defined against the Hive service, so the same policy covers a table in Hive and in Impala; Impala with Ranger covers policy design, column masking and row filtering. Keytabs and private keys should be readable only by the impala user, and rotating certificates needs a rolling restart, which the drain procedure above makes safe.

Monitoring

Watch three layers: daemon health, resource pressure and query outcomes. Daemon health is membership: the count of live backends seen by each coordinator should equal the number you expect. Resource pressure is memory in use against the limit, admission queue lengths and time spent queued, and scratch bytes written. Query outcomes are the rate of failed queries, rejected admissions and queries hitting their timeout. The metric names differ between releases, so rather than hard-coding them, scrape /metrics?json and select by pattern:

import json, re, sys, urllib.request

def metrics(host, port=25000):
    """Flatten an Impala daemon's /metrics?json page into {name: value}."""
    doc = json.load(urllib.request.urlopen(f"http://{host}:{port}/metrics?json", timeout=5))
    out = {}
    def walk(node):
        if isinstance(node, dict):
            if "name" in node and "value" in node:
                out[node["name"]] = node["value"]
            for v in node.values():
                walk(v)
        elif isinstance(node, list):
            for v in node:
                walk(v)
    walk(doc)
    return out

if __name__ == "__main__":
    host, pattern = sys.argv[1], re.compile(sys.argv[2])   # e.g. coord1 'admission|mem|queries'
    for name, value in sorted(metrics(host).items()):
        if pattern.search(name):
            print(f"{name}\t{value}")

Run it once against each daemon type to discover the names your release exposes, then feed the chosen ones into your metrics system. Archive query profiles from the coordinator's /queries page so a slow query can be compared with last week's run.

Worked example: the Monday-morning slowdown

Dashboards are slow every Monday from 08:00. The runbook goes layer by layer. /backends on both coordinators shows all 40 executors, so membership is fine. /admission shows the BI pool queueing for up to four minutes while the ETL pool is running three large queries. The profiles of those ETL queries show broadcast joins of a fact table that was loaded on Sunday night, and SHOW TABLE STATS confirms the new partitions have -1 rows.

The fix has two parts. Short-term, run COMPUTE INCREMENTAL STATS on the new partitions; the planner switches to partitioned joins and the ETL queries use a fraction of the memory. Long-term, add the stats statement to the end of the Sunday load job, and cap the ETL pool's memory so that BI queries always have admission headroom. No daemon failed, which is typical of Impala incidents.

Failure modes and trade-offs

  • Global INVALIDATE METADATA. Clears the whole catalog and causes a storm of reloads. Always name the table.
  • Mixed roles on large clusters. Coordinators that also execute run out of memory for both jobs, and every executor restart drops sessions. Split roles.
  • Memory limit too high. The kernel's OOM killer, not Impala, ends the process, and every query on the host fails. Leave headroom.
  • Scratch on the data disks. A large spill competes with scans for I/O and can fill disks. Use dedicated SSDs and per-directory caps.
  • Stale metadata after external writes. Event processing off, or writers that bypass the metastore, give wrong results rather than errors. Refresh as part of the writing job.
  • Single statestore or catalog. Without HA, a catalogd crash blocks all DDL and refreshes until it restarts and reloads. Enable HA where the release supports it.

What to do next

  1. Split coordinators and executors, and put at least two coordinators behind a load balancer with session affinity.
  2. Set an explicit memory limit and dedicated, capped scratch directories on every executor.
  3. Enable local catalog mode and HMS event polling, and remove routine INVALIDATE METADATA calls from jobs.
  4. Enable catalogd and statestored HA if your release supports them.
  5. Adopt the drain-and-restart script for every executor restart and upgrade.
  6. Schedule incremental stats after each partition load and alert on tables with no stats.
  7. Scrape the metrics endpoints for membership, memory, admission queues and failed queries, and archive query profiles.
Key takeaway: Impala administration is the care of three daemons and one shared metadata cache. Split coordinators from executors, give executors an explicit memory limit and dedicated scratch space, and use local catalog mode with HMS event polling so metadata stays current without blanket invalidation. Enable statestore and catalog HA where available, drain nodes with the shutdown command before restarting them, keep statistics fresh after every load, wire in Kerberos, TLS, LDAP and Ranger, and monitor membership, memory, admission queues and query failures from the daemons' own web endpoints.