Every Impala daemon, impalad, runs the same binary and can play two roles. As a coordinator it accepts client connections, plans queries, asks admission control for resources, hands work to other daemons, gathers results and returns them. As an executor it runs the fragments of a query plan: scans, joins, aggregations and exchanges. By default one daemon does both, and on small clusters that is fine. On larger ones, how you split the roles decides where memory goes, what a crash takes down and which component becomes the bottleneck.

This page follows a single query through both roles so that each responsibility is concrete, then covers sizing, failure semantics, query retries and the diagnostics that tell you which role is in trouble. For deployment and scaling of a whole cluster see Impala administration and scaling Impala; here the focus is the division of labour itself.

One binary, two roles

Two startup flags control the roles, and both default to true. --is_executor=false makes a daemon coordinator-only; --is_coordinator=false makes it executor-only. A cluster is healthy only if at least one daemon is a coordinator and at least one is an executor.

ResponsibilityCoordinatorExecutor
Client sessions (HiveServer2 protocol, impala-shell, JDBC, ODBC)YesNo
Metadata cache from catalogdYes, kept in the Java heapNot needed
Parse, analyse, plan (Java frontend)YesNo
Admission control decisionYesNo
Scan, join, aggregate, sort fragmentsOnly the root fragmentYes, almost all of the work
Result delivery and optional spoolingYesNo
Query profile aggregationYesReports its own instances

Both roles subscribe to the statestore for cluster membership, so every coordinator knows which executors are alive. Only coordinators need the catalog: the metadata about tables, partitions and files is used for planning, and an executor receives everything it needs to scan, down to file paths and byte ranges, inside the plan fragments the coordinator sends it.

One query through the coordinator

Take SELECT region, sum(amount) FROM sales WHERE day >= '2026-09-01' GROUP BY region against a partitioned Parquet table. The coordinator that owns the client session does the following, in order.

  1. Parse and analyse. The Java frontend resolves table and column names against the coordinator's metadata cache. If the table's metadata is not loaded, the coordinator waits for it, which is why the first query after a restart or invalidation is slow on the coordinator and not on executors.
  2. Plan. The planner prunes partitions to September onwards, lists the files and blocks that remain, and builds a distributed plan: a scan with a pre-aggregation on executors, a hash exchange on region, a merge aggregation, and a root fragment that gathers results.
  3. Admit. The coordinator estimates per-host memory and asks admission control whether the query fits in its resource pool. It runs, queues or is rejected; see admission control.
  4. Schedule. The scheduler assigns scan ranges to executors, preferring local or cached replicas, and decides how many fragment instances run on each host.
  5. Start. The coordinator sends each executor its fragment instances in one RPC per host, then starts the root fragment locally.
  6. Collect. Executors stream rows through exchanges, and the root fragment receives the final rows. The client fetches them, optionally from a spooled buffer.
  7. Finish. Executors report status and profile counters periodically and at completion; the coordinator merges them into one query profile, releases admitted resources and closes the query when the client does.
ClientJDBC / shellLoad balancersession affinityCoordinatorplan, admit, scheduleRoot fragmentfinal merge, resultscatalogdmetadata topicstatestoredmembershipExecutor 1scan + joinExecutor 2scan + joinExecutor 3agg + exchangeExecutor N...HDFS / S3 / Ozonedata filesSQLcatalogfragmentsrowsCoordinators hold sessions and metadata; executors hold data-processing memory and read the files.
Roles in a split cluster. Executors do not subscribe to catalog metadata; they receive scan ranges inside their fragments.

What executors do

On each executor the work is local and memory-bound. The executor creates fragment instances, starts scanner threads that read Parquet pages from storage, evaluates the date predicate, applies runtime filters arriving from other fragments, and pre-aggregates by region before hashing rows to the exchange. Another fragment on the same or other executors receives those partial sums and merges them.

Each executor enforces the query's per-host memory limit and its own process limit. Joins and aggregations that exceed their reservation spill to local scratch disk, if scratch directories are configured, rather than failing. None of this needs the catalog, which is why executors can run with a small Java heap and leave almost all memory to the backend. The internals are covered in query execution.

Why dedicated coordinators

When every daemon is both, three things go wrong as clusters grow. Every daemon keeps a full metadata cache, so a large catalog is paid for N times. Every daemon accepts connections, so a heavy client fetch or a big result set competes for memory with joins on the same host. And query fan-out is spread unpredictably, because whichever daemon a client lands on also coordinates.

Splitting fixes all three. The Impala documentation suggests a rough ratio of one coordinator for every 50 executors, one dedicated coordinator even on clusters under 10 nodes, a large JVM heap on coordinators for metadata, default JVM heaps on executors so memory goes to query processing, and peak resource use kept under 80 percent. Treat the ratio as a starting point: coordinator load scales with concurrent queries and fragment count, not with executor count alone.

Metadata on coordinators

Metadata is the largest and least predictable consumer of coordinator memory, so it deserves its own decision. In the legacy mode, catalogd broadcasts every table it has loaded through the statestore, and every coordinator keeps a full copy. A catalog with many partitions and small files makes each copy large, makes every change expensive to serialise, and makes a coordinator restart slow while it relearns everything.

On-demand metadata changes that. With --catalog_topic_mode=minimal on catalogd and --use_local_catalog=true on every coordinator, coordinators fetch metadata from catalogd at partition granularity when a query needs it and cache it locally, evicting under memory pressure. The cache size is set by local_catalog_cache_mb, which by default sizes itself to 60 percent of the Java heap, and entries expire after local_catalog_cache_expiration_s, one hour by default.

The trade is planning latency for memory. The first query to touch a cold table pays a fetch from catalogd, but adding a coordinator no longer multiplies the full catalog, and a restarted coordinator is ready quickly. On a cluster where coordinators are the scarce resource, on-demand mode is usually the cheapest single improvement, and it reinforces the split: executors never needed the catalog, and now coordinators hold only the part they use.

Worked example: laying out 64 hosts

Suppose a cluster of 64 hosts serves a BI tool with about 40 concurrent queries at peak and 15,000 partitions across its largest tables. A reasonable first layout is 2 coordinators and 62 executors, with the coordinators behind a load balancer. Two coordinators exceed the 1:50 rule, but the second is there for availability and rolling restarts, not capacity.

# coordinator flagfile (2 hosts)
--is_coordinator=true
--is_executor=false
--mem_limit=60%            # leaves room for a large JVM heap holding metadata

# executor flagfile (62 hosts)
--is_coordinator=false
--is_executor=true
--scratch_dirs=/data1/impala/scratch,/data2/impala/scratch

# client side, per session where results may be large or retries matter
SET SPOOL_QUERY_RESULTS=TRUE;
SET RETRY_FAILED_QUERIES=TRUE;

Set the coordinator heap through the daemon's JVM options in your distribution's configuration, sized from the catalog: measure heap use on a coordinator after a full metadata load and leave generous headroom. Configure the load balancer for session affinity on the HiveServer2 port, because Impala sessions and their open queries live on one coordinator; a client whose connections are spread across coordinators loses its session state.

Watch three numbers on the coordinators during the first week: JVM heap after garbage collection, the count of in-flight queries, and the time queries spend between submission and admission. Rising heap means metadata growth; rising in-flight counts with flat executor CPU mean clients are slow to fetch; long planning time points at the frontend. Any of these is a reason to add a coordinator; none of them is fixed by adding executors.

Executor groups

On larger or cloud deployments executors can be organised into executor groups. A query runs entirely within one healthy group, so a group is the unit of isolation and of scaling: adding a group adds concurrency without changing the plan for any single query. In Cloudera Data Warehouse, autoscaling adds and removes executor groups as queues build. On the coordinator side, --expected_executor_group_sets takes entries of the form prefix:size that describe sets of groups for workload-aware routing. The exact executor-side flags for joining a group vary by version, so check the documentation for the release you run before relying on them.

Groups change the failure arithmetic too. If one group drops below its required size it stops receiving queries while the others continue, which is a better outcome than every query on the cluster slowing down because a few hosts are missing.

Failure semantics by role

The roles fail differently, and the difference is the main operational reason to keep them apart.

EventEffectWhat helps
Executor crashes mid-queryEvery query with a fragment on it fails; the coordinator cancels the remaining fragmentsRETRY_FAILED_QUERIES; enough spare executors
Executor slow or unreachableCoordinator may blacklist it for subsequent schedulingInvestigate disk or network; blacklisting is temporary
Coordinator crashesAll its sessions and queries are lost; clients must reconnectTwo or more coordinators behind a load balancer; client reconnect logic
Coordinator heap pressureSlow planning, long pauses, failed metadata loadsBigger heap, fewer tiny files and partitions, on-demand metadata
Slow client fetchQuery holds executor memory until rows are consumedResult spooling; fetch timeouts; idle query timeouts
Statestore unavailableRunning queries continue; membership and metadata updates stallMonitor it; restart it; see the statestore details

Transparent retries need care. RETRY_FAILED_QUERIES was added in Impala 4.0 and is off by default in Apache Impala; some distributions enable it, so check yours. It retries the whole query, never individual fragments, applies only to SELECT statements, not INSERT or DDL, and is skipped once any rows have reached the client, unless results are spooled with spool_all_results_for_retries. Enable it for read-only BI workloads, where a retry is invisible and cheap, and leave it off for anything that writes.

Diagnosing which role is the bottleneck

When queries are slow, decide first which role is responsible, using the query profile. Planning time and admission wait are coordinator time; if they are large, the problem is metadata, frontend load or pool limits, not executors. Large differences between the fastest and slowest instance of the same fragment point at one executor: a hot disk, a skewed key or a remote read. A long gap between the last fragment finishing and the query closing usually means the client stopped fetching.

Cluster-wide, compare coordinator and executor health separately. Coordinators should show modest CPU, a stable heap after collection and steady admission; executors should show CPU and memory rising and falling with load. A coordinator at full CPU while executors idle is the textbook signature of an undersized coordinator tier.

What to do next

  1. List every impalad and record its two role flags; make sure no executor is accidentally accepting client connections.
  2. Run at least two coordinator-only daemons behind a load balancer with session affinity, and one per roughly 50 executors after that.
  3. Size coordinator JVM heap from measured post-load metadata use and leave executors at default heaps.
  4. Configure scratch directories on executors so large joins spill instead of failing.
  5. Turn on result spooling for BI clients and RETRY_FAILED_QUERIES for read-only pools only.
  6. Alert on coordinator heap after collection, admission wait and in-flight queries, separately from executor metrics.
  7. When a query is slow, read planning, admission and per-instance times in its profile before changing anything.
Key takeaway: Every impalad can coordinate and execute; past a handful of nodes, split them. Coordinators hold client sessions and the metadata cache, plan, admit, schedule and return results, so they need large JVM heaps and high availability behind a session-affine load balancer. Executors run scan, join and aggregation fragments and need memory and scratch disk. Plan roughly one coordinator per 50 executors, never fewer than two for availability, and diagnose slow queries by role using the profile. Enable RETRY_FAILED_QUERIES only for read-only workloads, with result spooling where rows stream early.