Impala internals
Anatomy of a distributed query.
This is the machinery I work on every day. Press run, and watch a SQL query fan out across an Impala cluster — the same six stages the engine walks through for every single query.
sales_report.sql
SELECT region, SUM(revenue) AS total
FROM sales
WHERE year = 2026
GROUP BY region
ORDER BY total DESC
LIMIT 10;
3 executors · MT_DOP=4
12 fragment instances
F01 · scan fragment · runs on all 3 executors
- 00:SCAN HDFS [sales] — 4.2B rows, column pruning on
- 01:AGGREGATE partial — SUM per region, per node
- 02:EXCHANGE [HASH(region)] — shuffle to coordinator
F00 · coordinator fragment · runs on impalad :21000
- 03:AGGREGATE merge — finalize one row per region
- 04:TOP-N ORDER BY total DESC LIMIT 10 — blocking
The distributed plan. Data flows top-down here; Impala’s EXPLAIN prints it bottom-up.
10 rows returned
| region | total |
|---|---|
| West | 48,211,903 |
| East | 44,702,118 |
| Central | 39,884,502 |
| South | 31,209,774 |
| North | 27,553,091 |
| South-East | 22,418,336 |
| North-West | 19,076,540 |
| Mid-West | 15,992,817 |
| South-West | 11,340,205 |
| North-East | 8,417,662 |
SUMMARY — per-operator timings
| Operator | #Hosts | #Inst | Avg Time | #Rows |
|---|---|---|---|---|
| 00:SCAN HDFS | 3 | 12 | 812ms | 4.2B |
| 01:AGGREGATE | 3 | 12 | 305ms | 1.1M |
| 02:EXCHANGE | 3 | 12 | 96ms | 48 |
| 03:AGGREGATE | 1 | 1 | 41ms | 12 |
| 04:TOP-N | 1 | 1 | 3ms | 10 |
Fetched 10 row(s) in 1.84s
Timings are illustrative — the six stages and the mechanics are how Impala really runs a query.