Every Impala query runs a plan chosen before the first row is read, and most slow Impala queries are slow because of that plan: the wrong join distribution, a scan that reads every partition, or an estimate so far off that the query is admitted with the wrong memory. Impala shows you the plan before you pay for it with EXPLAIN, and shows what really happened afterwards with SUMMARY and PROFILE. Learning to read the two side by side is the most useful skill an Impala user can have.
This article walks through EXPLAIN output node by node, covers the parts that decide performance, compares estimates with actuals, and fixes a bad plan. Plan listings below are illustrative: exact labels and detail lines vary between Impala versions, but the structure is the same.
Where a plan comes from
The coordinator that receives a query runs the planner, which is Java code in the Impala frontend. It parses and analyses the SQL against catalog metadata, then builds a single-node plan: a tree of operators that would compute the answer on one machine. It chooses join order and join sides using table and column statistics. Then it builds the distributed plan: it decides how each join and aggregation is spread across executors, inserts EXCHANGE nodes wherever rows must move between machines, and cuts the tree into fragments at those exchanges.
A fragment is a piece of the plan that runs without network hops inside it. The backend runs one or more instances of each fragment, normally one per executor that holds data for the scans in it (more when MT_DOP is raised). Exchanges connect instances: a sender streams row batches and a receiver on another fragment consumes them. So a plan tells you three things: what work is done, where it runs, and how much data crosses the network.
Anatomy of EXPLAIN output
Here is an illustrative plan for revenue per region over the last 30 days, joining a large sales fact table with a customers dimension:
Max Per-Host Resource Reservation: Memory=34.00MB Threads=5
Per-Host Resource Estimates: Memory=1.2GB
PLAN-ROOT SINK
|
08:MERGING-EXCHANGE [UNPARTITIONED]
| order by: sum(amount) DESC
| limit: 10
|
07:TOP-N [LIMIT=10]
| order by: sum(amount) DESC
|
05:AGGREGATE [FINALIZE]
| output: sum:merge(amount)
| group by: c.region
|
04:EXCHANGE [HASH(c.region)]
|
03:AGGREGATE [STREAMING]
| output: sum(s.amount)
| group by: c.region
|
02:HASH JOIN [INNER JOIN, BROADCAST]
| hash predicates: s.customer_id = c.id
| runtime filters: RF000 <- c.id
|
|--06:EXCHANGE [BROADCAST]
| |
| 01:SCAN HDFS [shop.customers c]
| partitions=1/1 files=8 size=410.00MB
| row-size=24B cardinality=12.00M
|
00:SCAN HDFS [shop.sales s]
partitions=30/1095 files=240 size=61.20GB
predicates: s.sale_date >= '2026-08-31'
runtime filters: RF000 -> s.customer_id
row-size=16B cardinality=2.10GThe header gives the per-host memory the query must reserve before it starts and the planner's estimate of what it will use; admission control uses these figures when no explicit memory limit is set. Each operator has a number (00, 01) that is its ID in every later report, a type, and detail lines. The tree is drawn with the main input straight below each operator and a second input on a |-- branch. For a hash join, the branch is the build side, loaded into a hash table first, and the straight line is the probe side streamed through it.
Reading bottom-up
The Impala documentation's advice is to read the plan from the bottom to the top, because that is the direction rows flow. For the plan above: node 00 scans only 30 of 1,095 partitions of sales thanks to the date predicate; node 01 scans all of customers; node 06 copies every customer row to every host running the scan of sales; node 02 builds a hash table on customer ID and probes it with sales rows, and publishes runtime filter RF000 back to scan 00 so rows with no matching customer are dropped at the scan; node 03 pre-aggregates on each host; node 04 shuffles partial sums by region so all partials for one region meet on one host; node 05 merges them; node 07 keeps the top ten on each host; node 08 merges those sorted lists on the coordinator; the root sink returns rows to the client.
Ask of every node: how many rows (cardinality), how wide (row-size), where it runs, and what crosses the network above it.
EXPLAIN_LEVEL: choosing the amount of detail
The EXPLAIN_LEVEL query option controls the output. Level 0 (MINIMAL) prints a flat list of operators, which is useful for a quick look at join order. Level 1 (STANDARD) is the default shown above. Level 2 (EXTENDED) adds more per-node detail, including memory estimates and the statistics the planner used, and is the level to use when you are checking whether statistics are present and sane. Level 3 (VERBOSE) groups the plan by fragment, which helps when you need to see exactly what runs where.
-- impala-shell
SET EXPLAIN_LEVEL=2;
EXPLAIN SELECT c.region, sum(s.amount) FROM shop.sales s
JOIN shop.customers c ON s.customer_id = c.id
WHERE s.sale_date >= '2026-08-31' GROUP BY c.region;
-- after a real run in the same session
SUMMARY;
PROFILE;Two warnings in the output deserve immediate action. A note that tables are missing relevant table or column statistics means cardinalities below are guesses. An estimate of -1 for rows or memory means the planner had no basis for a number at all.
Scans: partitions, files and filters
Scan nodes show partitions=used/total, the file count and total size. Partition pruning happens at planning time, so a scan reading 1,095 of 1,095 partitions when you expected 30 means the predicate could not be applied to the partition column: a function wrapped around the column, a type mismatch, or a filter on a different column. predicates lines list filters evaluated during the scan, and file formats such as Parquet can also skip row groups using min/max statistics.
runtime filters: RFnnn -> column on a scan means a join above will send a filter built from its build side, typically a Bloom filter or a min/max range, to skip non-matching rows or even whole files and partitions. The matching join lists RFnnn <- column. If the filter arrives late, the scan has already read data; the profile shows whether filters arrived and how much they eliminated.
Joins: broadcast versus partitioned
Every distributed hash join is either BROADCAST, where the whole build side is copied to every host running the probe side, or PARTITIONED (a shuffle join), where both sides are hash-partitioned on the join key so matching keys meet on the same host. The planner picks by estimated network and memory cost. Roughly, broadcast moves the build size times the number of probe hosts; partitioned moves the build size plus the probe size, once each.
With 20 hosts, broadcasting a 410 MB dimension moves about 8 GB and needs 410 MB of hash table on each host. Shuffling it together with 2 billion fact rows would move far more, so broadcast wins here. Now make the build side 40 GB: broadcasting it moves 800 GB and needs a 40 GB hash table on every host, which will spill or fail, so a partitioned join is the right choice. The planner can only make this call correctly if the size estimates are right, which is why statistics matter more than any hint.
The planner also orders joins and may swap sides so the smaller input becomes the build side; an INNER JOIN written one way can appear inverted in the plan. When the estimates are wrong and you cannot fix them yet, hints override the choice: JOIN /* +SHUFFLE */ big_table or /* +BROADCAST */ set distribution, and SELECT STRAIGHT_JOIN keeps the join order as written. Treat hints as temporary: they do not adapt when data grows.
Aggregation, sorting and exchanges
Grouped aggregates usually appear twice. A pre-aggregation on each host (sometimes labelled STREAMING) shrinks the data before it is sent, an EXCHANGE HASH(group keys) routes partials by key, and a FINALIZE aggregation merges them. If the pre-aggregation barely reduces rows, because the grouping key has nearly as many values as rows, Impala can pass rows through it rather than spend memory on it. TOP-N is a sort with a limit, kept in memory; a full SORT without a limit may spill to disk. ANALYTIC nodes implement window functions and are preceded by a sort on the partition and order keys.
Exchange types tell you the data movement: UNPARTITIONED gathers everything onto one instance (often the coordinator), HASH(...) shuffles by key, and BROADCAST copies to all. An unpartitioned exchange above a large input is a bottleneck, because one host must receive and process all of it.
Estimates versus actuals: SUMMARY
After a query runs, SUMMARY in impala-shell prints one line per plan node with columns including #Hosts, #Inst, Avg Time, Max Time, #Rows, Est. #Rows, Peak Mem and Est. Peak Mem. Two comparisons find most problems. #Rows against Est. #Rows shows where the planner's model broke: the first node, reading upward, where they differ by a large factor is where to look for missing or stale statistics or correlated predicates. Max Time against Avg Time shows skew: when one instance takes ten times longer than the average, one host received far more rows, usually because of a hot join or grouping key.
The script below reads a saved SUMMARY and flags both conditions. It parses the table that impala-shell prints, so check the column order against your version.
import re, sys
UNITS = {"K": 1e3, "M": 1e6, "B": 1e9}
def num(s):
s = s.strip()
if s in ("", "-1"):
return None # -1: planner had no estimate
m = re.fullmatch(r"([\d.]+)([KMB]?)", s)
return float(m.group(1)) * UNITS.get(m.group(2), 1) if m else None
def secs(s):
m = re.fullmatch(r"([\d.]+)(ns|us|ms|s|m)", s.strip())
if not m:
return 0.0
f = {"ns": 1e-9, "us": 1e-6, "ms": 1e-3, "s": 1, "m": 60}[m.group(2)]
return float(m.group(1)) * f
def check(path, ratio=10, skew=5):
for line in open(path):
cells = [c.strip() for c in line.strip().strip("|").split("|")]
# build-side rows are drawn as "|--06:EXCHANGE", which splits the operator cell
i = next((k for k, c in enumerate(cells) if re.match(r"-*\s*\d+:", c)), None)
if i is None or len(cells) < i + 7:
continue
op, avg, mx = cells[i].lstrip("- "), secs(cells[i + 3]), secs(cells[i + 4])
rows, est = num(cells[i + 5]), num(cells[i + 6])
if est is None:
print(f"{op}: no estimate (missing stats?)")
elif rows and max(rows, est) / max(min(rows, est), 1) >= ratio:
print(f"{op}: estimated {est:,.0f} rows, got {rows:,.0f}")
if avg > 0 and mx / avg >= skew:
print(f"{op}: max time {mx / avg:.1f}x average: skew")
check(sys.argv[1])
Worked example: a broadcast that should have been a shuffle
The numbers in this example are assumptions for illustration. A nightly report joins sales (about 2 billion rows in the 30-day window) with sessions, a new table of about 1.5 billion rows loaded last week. The query ran in four minutes on day one and now takes forty, sometimes failing with a memory-limit error. EXPLAIN shows HASH JOIN [INNER JOIN, BROADCAST] with sessions on the build side, a missing-statistics warning for sessions, and a build-side cardinality far below reality.
SUMMARY confirms it: the scan of sessions shows Est. #Rows of -1, the join's peak memory on every host is several times its estimate, and the join node's max time is close to its average, so this is not skew. With no statistics the planner judged the table small and chose to broadcast it to all 20 hosts, making every host build a hash table for the entire table. The fix is to run COMPUTE STATS shop.sessions (or COMPUTE INCREMENTAL STATS on a partitioned table that grows daily) and add it to the load job. The new plan uses PARTITIONED with EXCHANGE [HASH(...)] on both inputs, puts the smaller input on the build side, and the report returns to a few minutes. A /* +SHUFFLE */ hint would have fixed this one query too, but the statistics fix every query that touches the table.
Failure modes and trade-offs
- Missing or stale statistics. The most common root cause of bad join distribution and wrong admission memory. Compute stats as part of every load, not by hand.
- Estimates that are off for other reasons. Correlated predicates and skewed values mislead even good statistics. Compare with SUMMARY rather than trusting the plan.
- Hints that outlive the data. A broadcast hint written when a table was small becomes an outage when it grows. Review hints periodically.
- Unpruned scans. Functions on partition columns defeat pruning. Filter on the raw partition column.
- Single-host bottlenecks. Large unpartitioned exchanges, and skewed keys in partitioned joins, concentrate work on one host.
- Memory estimates that are too high. Queries then wait in admission queues while memory is actually free. Correct statistics, or set a sensible per-query memory limit.
What to do next
- Pick your three slowest recurring queries and run EXPLAIN on each at
EXPLAIN_LEVEL=2. - Read each plan bottom-up and write down partitions scanned, join distribution and build-side size.
- Run the query, save SUMMARY, and flag nodes where estimated and actual rows differ tenfold or max time is five times the average.
- Fix statistics first: add
COMPUTE STATSorCOMPUTE INCREMENTAL STATSto the load pipelines of every table involved. - Only then consider hints, and record why each one exists.
- Keep learning: how Impala executes fragments, runtime filters, admission control and memory estimates, catalog metadata and staleness and troubleshooting by symptom.