An Impala cluster looks like one SQL engine from the outside, but it is four kinds of process with very different jobs. There is the impalad daemon, which plans and runs queries. The statestored daemon tracks who is alive and passes cluster-wide state around. The catalogd daemon owns table metadata. And an optional admissiond daemon makes admission decisions for the whole cluster. Each holds different state and fails differently, and most production mistakes come from mixing them up.

This page catalogues the node types: what each process holds in memory, the flags that select its role, how it can be made highly available, and exactly what breaks when it dies. Flag names and defaults come from the Apache Impala source and documentation, read on 2026-10-03.

Four daemons, four kinds of state

Start with what each process remembers, because that decides everything else.

DaemonState it holdsWho talks to itIf it is lost
impalad (coordinator)Client sessions, query plans, a catalog cache (or on-demand metadata), query state and profilesClients; executors; statestored; catalogdQueries it was coordinating fail; clients must reconnect elsewhere
impalad (executor)Running fragment instances, scan buffers, hash tables, spill files, the data cacheCoordinators and other executors over KRPCQueries with fragments on it fail or are retried; capacity shrinks
statestoredSoft state only: the membership list and topic contents, rebuilt from subscribers after a restartEvery daemon, as subscribersRunning queries continue; new failures and metadata changes stop propagating
catalogdThe authoritative metadata cache loaded from the Hive Metastore and storageCoordinators (directly for DDL, via statestored for broadcasts)DDL and metadata loads fail; cached reads may continue
admissiondPool queues and admitted-resource accounting for the whole clusterCoordinators, through an RPC clientCoordinators that depend on it cannot admit new queries

Only coordinators and the catalog hold state that is costly to rebuild, and only coordinators face clients. That is why the usual production shape is a small number of large coordinators, a large number of executors, and a single (or HA-paired) instance of each control daemon.

statestoredmembership + topicscatalogd (active)metadata sourceadmissiond (optional)cluster-wide admissioncatalogd (standby)enable_catalogd_hastatestored (standby)enable_statestored_haimpalad coordinatoris_coordinator=true, is_executor=falseimpalad coordinatoris_coordinator=true, is_executor=falseimpalad executoris_coordinator=falseimpalad executoris_coordinator=falseimpalad executoris_coordinator=falseimpalad executor...heartbeats, topicscatalog updatesadmit / queuefragments (KRPC)Clients (impala-shell, JDBC, ODBC) connect only to coordinators. Executors never see a client session.
Node types in a production Impala cluster with dedicated coordinators, HA control daemons and an optional admission service.

impalad: one binary, two flags

One binary, impalad, plays both data-plane roles. Two boolean flags choose between them, and both default to true. In the Impala source, is_coordinator is described as: if true, this daemon can accept and coordinate queries from clients; if false, it will refuse client connections. is_executor is described as: if true, this daemon will execute query fragments.

is_coordinatoris_executorRoleWhen to use it
truetrueMixed (the default)Small clusters and development, where every host can plan and run work
truefalseDedicated coordinatorLarge or highly concurrent clusters; isolates planning, metadata and result handling from scan memory
falsetrueDedicated executorThe bulk of a large cluster; no client sessions and no planning, just fragments
falsefalseNeitherNo useful role. Treat it as a configuration error and alert on it

A coordinator plans SQL, keeps metadata for the tables it plans against and buffers results for clients, so its memory tracks metadata size and client behaviour. Executor memory tracks data volume and join sizes. Splitting the roles gives each a sizing model you can reason about.

Executors can be grouped. The executor_groups flag on an executor names the group it joins and can carry a minimum size, for example --executor_groups=default-pool-1:3. The flag's help text notes that a group is considered healthy for admission only once membership contains at least that many executors, and that currently only a single group may be specified per daemon.

Coordinators can also avoid holding the whole catalog. With --use_local_catalog=true on coordinators and --catalog_topic_mode=minimal on catalogd, a coordinator fetches metadata on demand from catalogd and caches only what it plans against, rather than receiving every catalog object through the statestore, which matters on clusters with many thousands of tables. For the per-query lifecycle on each role, see coordinator and executor roles in depth.

statestored: disposable by design

The statestore is a publish-subscribe hub. Every impalad and the catalogd register as subscribers, receive heartbeats, and receive updates to named topics: cluster membership, catalog updates and admission-control statistics among them. When a daemon stops answering heartbeats, the statestore removes it from membership and every coordinator stops scheduling work there.

Its state is soft. Nothing it holds is the only copy of anything. That is why Cloudera's documentation can say that if the statestore is not running or becomes unreachable, the Impala daemons continue running and distributing work among themselves as usual. The same documentation spells out the cost. The cluster becomes less robust if other daemons fail, because nobody is announcing the failure. Metadata becomes less consistent while the statestore is offline. And a DDL statement issued while it is down can leave queries that access the newly created object failing, because the change cannot be broadcast.

A statestore outage therefore degrades over time rather than stopping the cluster. For clusters that cannot tolerate even that, StateStore HA runs two statestored instances. Per the Apache Impala HA documentation, the statestored pair is started with enable_statestored_ha=true, state_store_ha_port (default 24020), state_store_peer_host and state_store_peer_ha_port. Every subscriber (catalogd, coordinators and executors) is restarted with enable_statestored_ha=true plus state_store_host, state_store_port, state_store_2_host and state_store_2_port. The standby takes over only when it has missed a run of heartbeat requests from the primary and a majority of clients have also lost their connection to the primary. That second condition stops a network split between the two statestores from producing two primaries. The topic mechanics are covered in the statestore architecture page.

catalogd: the single writer of metadata

catalogd is the single writer of metadata. It loads table definitions from the Hive Metastore and file and block information from the storage layer. It applies DDL issued through Impala, and it publishes changes. Coordinators receive those changes either as broadcast catalog-topic updates through the statestore or, in local-catalog mode, by fetching objects on demand. A coordinator executing CREATE TABLE, ALTER TABLE, REFRESH or INVALIDATE METADATA sends the operation to catalogd directly and waits for it.

Sizing catalogd is a metadata-volume problem, not a query-volume problem: a table with hundreds of thousands of partitions and millions of files costs heap to hold, time to load after an invalidate and bandwidth to propagate.

Catalog HA is active-passive. The Apache documentation says to set enable_catalogd_ha=true on both catalogd instances and on the statestore. The active statestore then assigns roles, designating one catalogd as active and the other as standby. If the active instance fails, the statestore promotes the standby and notifies all coordinators. The documentation is explicit that queries running at failover time can fail because they lose access to metadata, and must be rerun. HA shortens a catalog outage; it does not make failover invisible. Catalog and statestore HA are independent: you can enable either or both. Catalog HA depends on a working statestore to assign roles, so enabling it without statestore HA leaves the statestore as the remaining single point of failure. See the catalog and metadata plane for how updates flow.

admissiond: centralised admission

By default, admission control runs inside each coordinator. Each coordinator decides whether to admit, queue or reject a query against a resource pool's limits, using its own view plus pool statistics that other coordinators publish through the statestore. Because those statistics arrive with a delay, several coordinators admitting at once can each see room that the others are about to take. Limits are therefore approximate under bursts of concurrent submissions.

The admission control service moves that decision into one process. The Impala packaging script lists admissiond alongside impalad, catalogd and statestored as a startable service. On coordinators the switch is a single flag: in the Impala source, ExecEnv::AdmissionServiceEnabled() returns true exactly when admission_service_host is non-empty. In that case the coordinator uses a remote admission client, and otherwise a local one. The commit that introduced the RPC service lists the operations a coordinator calls: AdmitQuery, GetQueryStatus, ReleaseQueryBackends, ReleaseQuery and CancelAdmission.

The trade is the usual one for centralising a decision. You get one consistent view of pool usage and no over-admission races, which matters most when many coordinators share tight pools. You also get a new dependency in the admission path of every query. If admissiond is down, coordinators configured to use it cannot admit new work, whereas the distributed default keeps admitting with stale statistics. Pool configuration itself is covered in Impala admission control.

Ports and the debug surface

Defaults from the Apache Impala ports reference; each can be changed by the named flag. When a firewall rule is wrong, this table shows which role is cut off.

DaemonPurposeFlagDefault
impaladBeeswax clients (impala-shell)--beeswax_port21000
impaladHiveServer2 clients (JDBC/ODBC)--hs2_port21050
impaladHiveServer2 over HTTP--hs2_http_port28000
impaladKRPC between daemons--krpc_port27000
impaladStatestore subscriber--state_store_subscriber_port23000
impaladDebug web UI--webserver_port25000
statestoredSubscriber registration--state_store_port24000
statestoredDebug web UI--webserver_port25010
catalogdCatalog service--catalog_service_port26000
catalogdStatestore subscriber--state_store_subscriber_port23020
catalogdDebug web UI--webserver_port25020

Only coordinators need client ports exposed beyond the cluster. The statestore web UI shows which subscribers it believes are alive; the catalog UI shows which objects are loaded.

Worked example: a 40-host layout and a failover

Take a 40-host cluster serving dashboards (many short queries) and a nightly batch of heavy joins. A reasonable layout is two dedicated coordinators behind a load balancer, 36 executors, and two control hosts that each run a statestored and a catalogd as HA pairs. That keeps a single host failure from taking out both copies of any control daemon. The flag files might look like this:

# coordinators (2 hosts): plan, hold metadata on demand, face clients
--is_coordinator=true
--is_executor=false
--use_local_catalog=true
--enable_statestored_ha=true
--state_store_host=ctl-1
--state_store_port=24000
--state_store_2_host=ctl-2
--state_store_2_port=24000

# executors (36 hosts): fragments only, one named group
--is_coordinator=false
--is_executor=true
--executor_groups=default-pool-1:30
--enable_statestored_ha=true
--state_store_host=ctl-1
--state_store_port=24000
--state_store_2_host=ctl-2
--state_store_2_port=24000

# statestored on ctl-1 (mirror on ctl-2 with the peer pointing back)
--enable_statestored_ha=true
--enable_catalogd_ha=true
--state_store_ha_port=24020
--state_store_peer_host=ctl-2
--state_store_peer_ha_port=24020

# catalogd on ctl-1 and ctl-2
--catalog_topic_mode=minimal
--enable_catalogd_ha=true
--enable_statestored_ha=true
--state_store_host=ctl-1
--state_store_port=24000
--state_store_2_host=ctl-2
--state_store_2_port=24000

The minimum size of 30 on the executor group is a deliberate choice. With it, losing six executors still leaves the group healthy, but a rack outage that takes out more does not let queries run on a cluster too small to finish them in time. Use your Impala version's documentation to confirm that each flag above exists in your release before rolling it out. HA in particular arrived in Impala later than the core role flags.

Now walk through a failure. ctl-1 loses power at 02:10, during the batch. The standby statestored on ctl-2 sees missed heartbeats and finds that most subscribers have also lost ctl-1, so it takes over. Subscribers move to it. Because ctl-1 also hosted the active catalogd, the new primary statestore promotes the catalogd on ctl-2 and notifies coordinators. Running batch queries that needed metadata at that moment fail and must be rerun. Queries already executing on executors with their plans in hand are not affected by the catalog change itself. Without HA, DDL and failure detection would have stopped until ctl-1 returned.

Failure modes by node type

SymptomLikely node typeFirst check
New queries fail with metadata or catalog errors; running ones mostly finishcatalogdCatalog web UI reachable? Which instance is active?
DDL hangs or times outcatalogd, or the Hive Metastore behind itcatalogd logs for metastore calls; table load times
A dead host keeps receiving fragments and queries fail on itstatestoredStatestore web UI subscriber list; heartbeat errors
Newly created tables invisible to some coordinatorsstatestored (topic propagation) or local catalog cacheTopic versions on each coordinator; run REFRESH on the table
Clients cannot connect, but SQL runs fine on other hostsCoordinator, or the load balancer in front of itHS2 port reachability; is_coordinator on that host
Queries queue forever while the cluster is idleadmissiond, or executor group below minimum sizeAdmission service reachability; executor group membership count
Coordinator out of memory with modest query loadCoordinator holding a full catalogCatalog size; consider use_local_catalog

Two operational habits prevent most of these. First, alert per role, not per host: "executors below group minimum" and "no active catalogd" are the conditions that matter, and a host-level CPU alert says neither. Second, never colocate both members of an HA pair, or put a control daemon on a busy executor. The architecture-level picture of how these fit together is in Impala architecture.

Trade-offs

Dedicated coordinators cost hosts that do not scan data, which looks wasteful on a small cluster. Below roughly a dozen hosts, the mixed default is usually right. Local-catalog mode trades first-access latency (a coordinator must fetch metadata it has not seen) for a much smaller steady-state footprint. HA pairs add hosts and configuration in exchange for shorter, self-healing outages, and admissiond trades a new dependency for exact pool accounting.

What to do next

  1. For every host in your cluster, record its node type and the flags that make it so; flag any impalad with both is_coordinator and is_executor false.
  2. If you run more than about a dozen hosts with mixed roles, plan a move to dedicated coordinators and measure coordinator heap before and after enabling use_local_catalog.
  3. Decide whether catalogd and statestored outages are acceptable; if not, design HA pairs on separate hosts and rehearse a failover in staging, rerunning the queries it fails.
  4. Set an executor group minimum size from your heaviest regular query, and alert when membership falls below it.
  5. If concurrent coordinators over-admit into tight pools, evaluate admissiond and add it to your control-plane monitoring before enabling it.
  6. Verify every flag on this page against your Impala release's documentation before changing production configuration.
Key takeaway: Impala has four node types. impalad, with is_coordinator and is_executor deciding whether it plans, executes or both. statestored, which holds only soft state and can be lost briefly without stopping running queries. catalogd, the single writer of metadata, whose loss stops DDL and metadata loads. And the optional admissiond, which centralises admission at the cost of a new dependency. Size coordinators for metadata and clients, executors for data, and use the HA pairs and executor group minimums so that a single host failure produces a short, explainable degradation instead of an outage.