Most slow or failed Impala queries that involve joins come down to one decision made at planning time: how each join's inputs are distributed across the executors. Send the smaller table to every node and the join is a cheap local lookup. Send a large table to every node and the query floods the network, runs out of memory, or spills to disk. Statistics drive that decision, and when they are missing Impala falls back to a default that is often wrong for large tables.
This article explains the two join operators, the two distribution strategies and what each costs, how the planner orders tables and chooses a strategy, which query options and hints change its mind, and how to confirm what happened. A worked example puts numbers on the trade-off, and the article ends with a checklist for your own queries. For the equivalent decisions in Hive see Hive join strategies.
Two operators, two distributions
Impala joins with a hash join whenever there is at least one equality predicate between the inputs. It reads the right-hand input, the build side, into an in-memory hash table keyed on the join columns, then streams the left-hand input, the probe side, and looks up each row. Build cost is paid in memory; probe cost is paid in CPU and is proportional to the probe rows.
Without an equality predicate, for example a range condition or a cross join, there is nothing to hash on, and Impala uses a nested loop join, which compares rows pairwise. It is correct for any condition and expensive for large inputs, so treat one in a plan over large tables as a warning.
Independently of the operator, the planner chooses a distribution for a hash join. Broadcast sends the entire build side to every node that runs the probe side; each node builds a full copy of the hash table and probes it with the rows it scanned locally. Partitioned, which hints and options call SHUFFLE, hashes both sides on the join key and sends each row to the node that owns its hash bucket, so each node builds and probes only its share.
What each strategy costs
For broadcast, network traffic is roughly the build side's size times the number of nodes receiving it, and every one of those nodes must hold the whole hash table in memory. The probe side does not move at all.
For partitioned, traffic is roughly the build side plus the probe side, each sent once, and each node holds about 1/N of the hash table. The probe side, usually the big fact table, now crosses the network, which is the price.
The planner compares estimates of this kind using row counts and average row sizes from table and column statistics, after applying predicate selectivity. When the broadcast estimate is lower it broadcasts. Two guards sit on top of that comparison. The BROADCAST_BYTES_LIMIT query option, 32 GB by default, caps the estimated size of a broadcast input; above it the planner chooses a partitioned join even if the cost comparison favoured broadcast, and 0 disables the cap. And joins whose semantics need a particular distribution, described below, are not offered the other one.
Join order and the build side
With statistics available for every table, Impala reorders the joins in a query itself. The shape it aims for, which Impala's documentation also recommends when you order tables by hand, is the largest table first on the left, so it is scanned and streamed rather than held in memory, then the smallest table, so intermediate results shrink early, then progressively larger tables. Every right-hand input becomes a build side, so this ordering keeps hash tables small.
Without statistics the ordering degrades in a specific way. Tables that have statistics are placed on the left in descending order of estimated cost; tables without statistics are treated as zero-size and always placed on the right. A huge fact table with no statistics therefore becomes a build side, and being estimated at zero bytes it looks ideal to broadcast.
Whether a missing-statistics join is broadcast or partitioned is controlled by DEFAULT_JOIN_DISTRIBUTION_MODE, which takes BROADCAST, the default, or SHUFFLE. Impala's documentation recommends SHUFFLE when setting up new clusters: an unnecessary shuffle of a small table costs a little, while an unnecessary broadcast of a large one can take down the query. The real fix is statistics; statistics and metadata covers keeping them current.
-- Safer default for a cluster where some tables may lack statistics
SET DEFAULT_JOIN_DISTRIBUTION_MODE=SHUFFLE;
-- Gather statistics so the planner can choose for itself
COMPUTE STATS sales.orders;
COMPUTE INCREMENTAL STATS sales.order_lines PARTITION (dt='2026-10-02');
Join types that constrain distribution
Some join semantics cannot be computed correctly under one of the strategies, whatever the sizes. Reasoning them through from first principles tells you what you will see in a plan.
- Right outer and full outer joins must emit each unmatched build row exactly once. With broadcast every node holds every build row and none of them knows whether another node matched it, so these joins need partitioned distribution, unless the planner can swap the inputs and turn the join into a left outer join.
- NOT IN subqueries become null-aware anti joins: a single NULL in the subquery result changes the answer for every probe row, so every node must see the whole build side, which is a broadcast.
- Nested loop joins have no key to hash rows on, so they cannot be partitioned; their build side is broadcast. A non-equi join against a large table is expensive for this reason too.
- Inner and left outer joins work either way, and that is where the cost comparison and your hints apply.
Reading the choice in EXPLAIN
The distribution is printed on every join node, and the exchange beneath the build side shows how its rows travel. A trimmed plan:
04:HASH JOIN [INNER JOIN, BROADCAST]
| hash predicates: o.customer_id = c.id
| runtime filters: RF000 <- c.id
|
|--03:EXCHANGE [BROADCAST]
| |
| 01:SCAN HDFS [sales.customers c]
| predicates: c.country = 'DE'
|
00:SCAN HDFS [sales.orders o]
runtime filters: RF000 -> o.customer_idA partitioned join shows [INNER JOIN, PARTITIONED] with an EXCHANGE [HASH(...)] on both inputs. Check the cardinality and memory estimates at a higher explain level, and after running compare estimated against actual rows in the summary; reading Impala query plans walks through both.
Worked example: one join, three answers
A cluster has 20 executors. orders holds 2 TB of scanned data; customers holds 40 GB. The query joins them on customer ID.
Unfiltered, broadcasting customers sends about 40 GB x 20 = 800 GB and asks every node to hold a 40 GB hash table plus overhead. Partitioning sends about 2 TB + 40 GB once and each node holds about 2 GB of hash table. Broadcast moves less data on paper, but its estimated input exceeds the 32 GB BROADCAST_BYTES_LIMIT, so the planner partitions, and that is the right call: a 40 GB hash table on every node would exhaust per-node memory limits.
Add WHERE c.country = 'DE' and suppose statistics estimate 10 percent of customers survive. The build side is now about 4 GB: broadcast costs about 80 GB of traffic against more than 2 TB for partitioned, and a 4 GB table per node is affordable. The planner broadcasts.
Now drop the statistics on customers. It is treated as zero-size, placed on the right, and with the default mode broadcast. In the filtered query that happens to be fine. In the unfiltered one there is no usable size estimate to compare against BROADCAST_BYTES_LIMIT, so the default mode decides and 40 GB goes to every node. Same SQL, a missing COMPUTE STATS, and a query that previously ran now fails for lack of memory.
Hints: when and how
Hints override the distribution of one join. Place /* +BROADCAST */ or /* +SHUFFLE */ immediately after the JOIN keyword; it applies to the tables on either side of that join. Because the planner may reorder joins, put STRAIGHT_JOIN straight after SELECT so the written order is kept and the hint lands where you meant it.
SELECT STRAIGHT_JOIN c.segment, SUM(o.amount)
FROM sales.orders o
JOIN /* +SHUFFLE */ sales.customers c ON o.customer_id = c.id
JOIN /* +BROADCAST */ ref.dates d ON o.order_date = d.dt
WHERE d.fiscal_quarter = '2026Q3'
GROUP BY c.segment;Hints cannot be applied to joins Impala creates internally from subqueries in a WHERE clause, so rewrite such a subquery as an explicit join if you need to steer it. Treat every hint as technical debt: it freezes a decision that was right for today's data. Fix statistics first, and hint only when the planner is still wrong with accurate statistics, with a comment saying why.
Memory, skew, spilling and runtime filters
The build side's hash table must fit within the query's per-node memory. When it does not, the join spills partitions of the hash table to scratch disk and processes them in passes, which is slower but completes; spill to disk explains how to size scratch space and what spilling costs.
Partitioned joins have a weakness broadcast does not: key skew. Every row for a given key goes to the same node, so one very common customer ID or a flood of NULL keys makes one node do most of the work and spill while the others wait. Filter NULL keys before the join when they cannot match, and look for one fragment instance with far more rows than its peers in the profile.
Hash joins also produce runtime filters. The build side, once complete, publishes a Bloom or min-max filter on its join keys that probe-side scans use to skip rows, row groups and whole partitions. A selective build side can therefore shrink the probe scan dramatically; runtime filters covers their types, timing and tuning.
Failure modes
| Symptom | Likely cause | Fix |
|---|---|---|
| Query fails with memory limit exceeded on a join | Large build side broadcast, often from missing stats | COMPUTE STATS; SHUFFLE default mode; SHUFFLE hint |
| One fragment runs far longer than the rest | Key skew in a partitioned join | Filter NULL keys; isolate hot keys |
| Huge network time, little CPU | Broadcast to many nodes of a mid-sized table | Check estimated build size; consider SHUFFLE |
| Plan changed overnight | Statistics refreshed or dropped; data grew past the broadcast cap | Compare EXPLAIN before and after |
| NESTED LOOP JOIN over large tables | Missing equality predicate | Add an equi-join key or pre-bucket the range |
What to do next
- List the tables in your heaviest queries and confirm each has current table and column statistics.
- Set DEFAULT_JOIN_DISTRIBUTION_MODE=SHUFFLE as the cluster default if any tables may lack statistics.
- Run EXPLAIN on the top ten queries and record each join's distribution and estimated build size.
- Compare estimated with actual rows in the profile; stale estimates are the root of most bad joins.
- Remove hints that were added for old data volumes, and comment the ones you keep.
- Check the profile of slow partitioned joins for skewed fragment instances.
- Confirm runtime filters are produced and applied for selective dimension joins.