Impala failures look varied, but they fall into a small number of classes, and each class leaves a characteristic trace. A query that was rejected never ran, so the evidence is in admission control. A query that cannot see a table is a metadata problem. A query that died with a memory error, or ran for ten minutes instead of ten seconds, left a runtime profile that says exactly which operator on which host did it.
This article is a method rather than a list of error messages. It sorts problems by symptom, shows how to capture and read the evidence, works one slow query from profile to fix, and gives a playbook per failure class with the settings that matter. Deep explanations of each subsystem live in their own articles and are linked where they apply; the goal here is to get you from something is wrong to this is the layer, and this is the change quickly and without guessing.
Start with the symptom, not the fix
The most common troubleshooting mistake is to reach for a setting (raise the memory limit, add a hint, restart the catalog) before knowing which stage failed. Every Impala query moves through the same stages: the coordinator parses and analyses it against cached metadata, plans it, asks admission control for resources, then runs fragments on executors that scan, join, aggregate and exchange data. A failure belongs to exactly one stage, and each stage has different evidence.
| Symptom | Stage | First evidence |
|---|---|---|
| Query rejected immediately or waits, then times out in the queue | admission control | error text naming the pool; the pool's page in the web UI |
| AnalysisException: table, column or partition not found | analysis (metadata) | compare Hive Metastore with Impala's cached view |
| Memory limit exceeded, row size or scratch errors mid-query | execution | profile: which fragment on which host hit the limit |
| Finished, but much slower than usual | planning or execution | ExecSummary: estimated vs actual rows, max vs average time |
| Finished with unexpected results after new data landed | metadata | was the table refreshed after the write? |
Capture the evidence first
In impala-shell, SUMMARY prints the ExecSummary of the last query and PROFILE prints the full runtime profile. The coordinator's debug web UI, on port 25000 by default, lists in-flight and recent queries with their plan, summary, profile and timeline; the statestore and catalog daemons serve their own pages on 25010 and 25020. Save the profile text to a file before anyone reruns the query with different settings, because a rerun replaces the evidence of the original behaviour.
-- impala-shell: run, then inspect without re-running
[coord:21050] > SET MEM_LIMIT=8g;
[coord:21050] > SELECT c.region, sum(o.amount)
> FROM orders o JOIN customers c ON o.cust_id = c.id
> WHERE o.dt >= '2026-09-01'
> GROUP BY c.region;
[coord:21050] > SUMMARY; -- ExecSummary: per-operator time, rows, memory
[coord:21050] > PROFILE; -- full runtime profile (save it)
-- Look at the plan before running
[coord:21050] > SET EXPLAIN_LEVEL=2;
[coord:21050] > EXPLAIN SELECT ...;The profile has three layers. The top has the query text, options in effect, the plan and a query timeline of milestones such as planning finished, admitted, first row fetched and unregistered. Then comes the ExecSummary table. Below that are per-fragment, per-instance counters for every operator. Read the timeline first: if most of the wall-clock time sits between submit for admission and admitted, the query was queued, not slow, and no amount of plan tuning will help. If it sits after first row fetched, the client was slow to fetch results.
Reading the ExecSummary
The ExecSummary is one row per plan operator, read bottom-up: scans at the bottom feed joins and aggregations that feed exchanges up to the coordinator. The columns that matter most are Avg Time vs Max Time (a big gap means one host did much more work: skew or a slow node), #Rows vs Est. #Rows (a big gap means the planner worked from wrong or missing statistics) and Peak Mem vs Est. Peak Mem (a big gap explains memory errors and bad admission decisions).
Operator #Hosts Avg Time Max Time #Rows Est. #Rows Peak Mem Est. Peak Mem Detail
-----------------------------------------------------------------------------------------------------
06:EXCHANGE 1 0.1ms 0.1ms 5 5 16 KB 16 KB UNPARTITIONED
05:AGGREGATE 10 2.1ms 3.0ms 50 50 2 MB 10 MB FINALIZE
04:EXCHANGE 10 0.3ms 0.5ms 50 50 64 KB 64 KB HASH(c.region)
03:AGGREGATE 10 110ms 950ms 50 50 2 MB 10 MB STREAMING
02:HASH JOIN 10 1.4s 11.8s 412M 1.2M 1.9 GB 34 MB INNER JOIN, BROADCAST
|--01:SCAN HDFS 10 8ms 12ms 3.1M -1 24 MB 32 MB customers c
00:SCAN HDFS 10 2.2s 2.6s 412M 1.2M 160 MB 88 MB orders oThis is the summary of the worked example below. Read it bottom-up. The scan of orders returned 412 million rows against an estimate of 1.2 million. The customers scan shows an estimate of -1, which means the table has no statistics at all. The planner therefore believed both sides were small and chose a broadcast join, sending the whole right side to every host. The join's max time is eight times its average, and its peak memory is 1.9 GB against an estimate of 34 MB. Three signals, one cause: missing statistics.
Worked example: from slow query to fix
The query joins a week of orders to customers and normally takes about ten seconds. After a reload of the customers table from Spark, it takes three minutes and sometimes fails with a memory error on one host. The profile timeline shows admission was instant, so this is an execution problem. The ExecSummary above shows the estimate gaps, and EXPLAIN confirms it by printing a warning that tables are missing relevant table and/or column statistics.
The reload replaced the customers table, which dropped its statistics, and the recent orders partitions had never had incremental stats computed. With no row counts the planner guessed, chose a broadcast join, and the one host holding the most matching orders rows also did the most probing, hence the max-versus-average gap. The fix is to give the planner real numbers, then verify the plan, and only reach for a hint if the planner is still wrong with correct statistics.
-- 1. Give the planner real numbers
COMPUTE STATS customers;
COMPUTE INCREMENTAL STATS orders PARTITION (dt='2026-09-29');
-- 2. Re-check the plan: estimates should now be close to actuals,
-- and the join strategy should reflect the true sizes
EXPLAIN SELECT ...;
-- 3. Only if the planner is still wrong, pin the strategy for this query
SELECT c.region, sum(o.amount)
FROM orders o JOIN /* +SHUFFLE */ customers c ON o.cust_id = c.id
WHERE o.dt >= '2026-09-01'
GROUP BY c.region;After COMPUTE STATS the estimates in EXPLAIN matched the actuals to within a few percent and the plan picked the strategy that fits the real sizes. Runtime returned to about ten seconds with no hint needed. The lasting fix was operational: the Spark job that reloads customers now runs COMPUTE STATS through Impala as its last step, and new partitions of orders get incremental stats when they land. How statistics are stored and when incremental stats help is covered in Impala statistics and metadata.
Playbook: rejected or queued queries
Admission control decides whether a query may start, based on the memory it is expected to need and the limits of its resource pool. Three messages matter. A rejection says the query's memory requirement is larger than the pool, or a per-host limit, could ever allow; waiting will not help, so lower the query's needs (better statistics, a smaller MEM_LIMIT) or route it to a larger pool. A queue timeout says the query waited longer than the pool's queue timeout, which defaults to 60 seconds, because other queries held the memory; this is contention, and the fix is capacity, scheduling or pool separation, not the query. A queue-full rejection says too many queries were already waiting.
Estimates drive admission, so missing statistics hurt twice: they cause bad plans and bad memory estimates, which either over-reserve and queue everything else or under-reserve and fail at run time. Setting MEM_LIMIT per query or a pool's minimum and maximum query memory limits makes admission use a bounded figure instead of the estimate. Pool design, queue sizing and the interaction with memory estimates are covered in Impala admission control.
Playbook: memory limit exceeded and spilling
A memory-limit error names the limit that was hit: the query's own limit on one host, the pool's, or the whole process on that daemon. The profile shows which fragment instance was using the memory, and that is the operator to look at. Hash joins, aggregations, sorts and analytic functions can spill to scratch disk when their memory is exhausted, so a memory error on one of them usually means spilling was not possible or not enough: scratch directories missing or full, SCRATCH_LIMIT reached, or a single hash table partition too large to process even after repartitioning, which happens with extreme key skew.
Spilling itself is not a failure, but heavy spill turns a memory problem into an I/O problem; look for large spilled-partition counts and scratch write volume in the operator counters. The usual fixes, in order: correct statistics so the planner sizes joins correctly, choose the build side and join strategy correctly, remove skew (filter or salt the hot key), then raise limits. A different error, the one saying a row could not be materialised and suggesting MAX_ROW_SIZE, is about a single very wide row, typically a huge string or nested value, and needs that option or a schema change rather than more memory. See Impala memory limits and spill to disk.
Playbook: stale or missing metadata
Impala caches metadata from the Hive Metastore and file listings from storage, and the catalog daemon distributes changes to coordinators. Changes made through Impala update the cache automatically. Changes made elsewhere do not, which produces two classic symptoms: a table created by Hive or Spark is reported as not found, and new files written into an existing table are invisible, so counts are too low.
-- A table or partition was created outside Impala (Hive, Spark, a new directory)
INVALIDATE METADATA sales.orders; -- reload this table's metadata from scratch
-- Files were added or replaced in an existing table/partition
REFRESH sales.orders; -- reload file and block metadata
REFRESH sales.orders PARTITION (dt='2026-09-30'); -- narrower and cheaper
-- Avoid: bare INVALIDATE METADATA with no table name.
-- It discards cached metadata for every table in the catalog.Use the narrowest statement that fixes the problem. REFRESH on a partition is cheap; INVALIDATE METADATA on a table forces a full reload on next access; a bare INVALIDATE METADATA flushes the whole catalog and causes a burst of slow first queries across every workload. If writers are yours, make the refresh part of the write job. Catalog daemon memory pressure, a very large catalog topic and slow metadata loading for tables with many partitions and files are separate, larger problems, discussed with local catalog mode and automatic invalidation in Impala metadata in depth.
Playbook: slow queries that finish
- Estimates far from actuals: compute statistics, as in the worked example. Nothing else in this list is worth trying until estimates are sane.
- Max time far above average on one operator: skew or one slow host. If the same host is slow for every operator, suspect the node (disk, noisy neighbour, remote reads); if one operator only, suspect key skew in that join or aggregation.
- Scans read far more data than expected: check that partition predicates actually prune (functions on the partition column can defeat pruning), that files are a columnar format, and that the table is not made of thousands of tiny files, each of which costs scheduling and open time.
- Probe side not filtered: runtime filters built from the join's build side should shrink the probe scan. If the filter arrived too late, the scan ran unfiltered;
RUNTIME_FILTER_WAIT_TIME_MScontrols how long scans wait. See runtime filters. - Time before the first row, not during execution: long planning usually means cold or huge metadata; long admission means queueing.
- Time after the first row: the client fetches slowly or the result is enormous. Fetching millions of rows through a BI tool is a design problem, not an Impala one.
Timeouts and cancelled queries
Queries that disappear are often cancelled by a timeout rather than by a fault. EXEC_TIME_LIMIT_S cancels a query that runs too long. QUERY_TIMEOUT_S cancels a query that has been idle, for example a client that stopped fetching, and the error says the query expired due to client inactivity. Idle session timeouts close whole sessions. Coordinator restarts and executor failures also cancel queries; the coordinator's log and the profile's error section name the host. Check which of these fired before tuning anything.
Operational habits that shorten incidents
- Archive profiles for failed and slow queries automatically.
- Alert on per-pool queue length and queue timeouts, daemon memory, catalog heap and scratch disk usage.
- Change one thing at a time and compare profiles before and after.
What to do next
- Pick the slowest recurring query on your cluster, run it in impala-shell and save its
SUMMARYandPROFILE. - Read the ExecSummary bottom-up and mark every operator where actual rows or memory are more than ten times the estimate.
- Run
EXPLAINand fix any missing-statistics warning withCOMPUTE STATSorCOMPUTE INCREMENTAL STATS. - List every job that writes to Impala tables from outside Impala and add the narrowest
REFRESHorINVALIDATE METADATAto it. - Review pool queue timeouts and memory limits against the last week's rejections and queue timeouts.
- Write the first three entries of your runbook: symptom, stage, evidence, fix.