Join strategies
Broadcast hash, shuffle hash, sort-merge and nested loop joins: how Spark picks one, and how to steer it.
On this page
Show code in
Every code block on the page follows this.
BETWEEN join can run for hours.Join types vs join strategies
| Join type | Returns |
|---|---|
inner | Only rows with a match on both sides |
left / right outer | Every row of one side, with nulls where the other side has no match |
full outer | Every row of both sides, matched where possible |
left_semi | Left rows that have at least one match; left columns only, never duplicated |
left_anti | Left rows that have no match |
cross | Every combination of left and right rows |
The type is part of what your query means. The strategy is an execution choice Spark makes during physical planning: the same inner join can run as any of the five strategies depending on table sizes, the join condition, configs and hints. The PySpark join and left_semi / left_anti lessons cover the types; this lesson covers the strategies.
The five strategies
| Strategy (plan node) | How it works | Needs | Main risk |
|---|---|---|---|
Broadcastbroadcast join: A join where the small table is copied to every machine, so the big table never has to be shuffled. Learn more → hash joinBroadcastHashJoin | Copy the small side to every executorexecutor: A worker process on a cluster machine that runs tasks and holds data in memory. Learn more →, build a hash table from it, stream the large side through it. No shuffleshuffle: Moving rows between machines so that all rows with the same key end up together. Needed by joins, groupBy and sorting, and usually the most expensive step of a job. Learn more →. | Equi-join; one side small; not a full outer join | Driverdriver: The process that runs your program, plans the work and sends tasks to executors. Learn more → and executor memory when the "small" side is not small |
Shuffle hash joinShuffledHashJoin | Shuffle both sides by key. In each partitionpartition: A chunk of a DataFrame's rows. Spark processes each partition as one task, so partitions decide how much work runs in parallel. Learn more →, build a hash table from the smaller side and probe it with the other. No sort. | Equi-join; the build side of every partition fits in memory | Out-of-memory on one large or skewedskew: When one key has far more rows than the others, so one task does most of the work while the rest wait. Learn more → partition |
Sort-merge joinSortMergeJoin | Shuffle both sides by key, sort each partition, then walk the two sorted streams together. | Equi-join; sortable key types | Two shuffles and two sorts; skewed keys |
Broadcast nested loop joinBroadcastNestedLoopJoin | Broadcast one side and compare every row of the other side with every broadcast row. | Any condition, or none; any join type | rows × rows comparisons; slow when both sides are large |
Cartesian productCartesianProduct | Pair every partition of one side with every partition of the other, and test the condition on every row pair. | Inner or cross join | Work and output grow as rows × rows |
An equi-join has at least one equality between a column of each side in its condition (o.customer_id = c.customer_id, possibly ANDed with other conditions). Only equi-joins can use the hash and sort strategies, because those need a key to partition and match on. A condition with no such equality, like e.ts BETWEEN s.start AND s.end, a.x < b.y, or two equalities joined by OR, is a non-equi join and can only run as a nested loop.
How Spark picks one
For an equi-join, Spark walks this list and takes the first strategy that applies:
- 1Hints. If you hinted a strategy and the join supports it, Spark uses it. When both sides carry different hints, the priority is BROADCAST, then MERGE, then SHUFFLE_HASH, then SHUFFLE_REPLICATE_NL.
- 2Broadcast hash join if the join type allows broadcasting a side and that side's estimated size is below
spark.sql.autoBroadcastJoinThreshold(10 MB). If both sides qualify, the smaller is broadcast. - 3Shuffle hash join if
spark.sql.join.preferSortMergeJoinis false (it is true by default), one side is small enough to hash per partition (below the broadcast threshold × the number of shuffle partitions), and it is at least three times smaller than the other side. - 4Sort-merge join if the join keys are sortable. Most joins of two large tables end here.
- 5Cartesian product if it is an inner or cross join.
- 6Broadcast nested loop join as the last resort. It handles every case, even if that means broadcasting a large side.
For a non-equi join the list is shorter: a broadcast hint or a side below the threshold gives a broadcast nested loop join; otherwise an inner or cross join becomes a Cartesian product; otherwise Spark falls back to a broadcast nested loop join anyway.
spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold (0 by default, so off).Why sort-merge is the default
A shuffle hash join skips the sort, so on paper it is faster. But its hash table for a partition must fit in memory, and one oversized or skewed partition fails the tasktask: The work for one partition in one stage, run on one CPU core. Learn more →. A sort-merge join sorts with an external sorter that spillsspill: Writing data to local disk because it does not fit in memory. Slower, but the job keeps running. Learn more → to disk, and the merge holds only the rows of the current key, so it slows down under pressure instead of failing. That robustness is why Spark prefers it. Shuffle hash join is worth trying when one side is too big to broadcast but small per partition, and keys are not skewed.
Which strategies support which join types
| Join type | Broadcast hash | Shuffle hash | Sort-merge | Broadcast nested loop | Cartesian |
|---|---|---|---|---|---|
| inner, cross | Yes, either side | Yes | Yes | Yes, either side | Yes |
| left outer, left semi, left anti | Right side broadcast | Yes | Yes | Right side broadcast | No |
| right outer | Left side broadcast | Yes | Yes | Left side broadcast | No |
| full outer | No | Yes (Spark 3.1+) | Yes | Yes, but scans data several times | No |
The pattern: the side that must keep unmatched rows is streamed, so only the other side can be the broadcast or build side.
Hints
orders.join(customers.hint("broadcast"), "customer_id") # BroadcastHashJoin orders.join(customers.hint("merge"), "customer_id") # SortMergeJoin orders.join(customers.hint("shuffle_hash"), "customer_id") # ShuffledHashJoin events.join(sessions.hint("shuffle_replicate_nl"), events.ts.between(sessions.start, sessions.end)) # CartesianProduct
SELECT /*+ BROADCAST(c) */ * FROM orders o JOIN customers c ON o.customer_id = c.customer_id; SELECT /*+ MERGE(c) */ * FROM orders o JOIN customers c ON o.customer_id = c.customer_id; SELECT /*+ SHUFFLE_HASH(c) */ * FROM orders o JOIN customers c ON o.customer_id = c.customer_id; SELECT /*+ SHUFFLE_REPLICATE_NL(s) */ * FROM events e JOIN sessions s ON e.ts BETWEEN s.start AND s.end;
A hint is a request, not an order. Spark ignores a hint it cannot honour, such as broadcasting the preserved side of an outer join or a shuffle hash join on a non-equi condition, and logs a warning. MERGEJOIN and SHUFFLE_MERGE are aliases of MERGE.
Making a slow non-equi join fast
- OR across two keys —
a.id = b.id OR a.email = b.emailhas no single equality, so it becomes a nested loop. Rewrite it as two equi-joins and union the results (deduplicating if a row can match both ways). - Range joins — Matching an IP to an IP range, or an event to a time window, has no equality. Add one the condition already implies: a coarse bucket (the /16 prefix of the IP, the hour of the timestamp). Explode each range into the buckets it covers, join on bucket equality, and keep the range test as an extra condition. The join becomes a sort-merge join on the bucket.
- Functions on the key — A join on
lower(a.code) = b.codeis still an equi-join, but nothing about it can use bucketing or existing partitioning. Normalise the key once, before the join.
Check what Spark chose
Search the physical plan for the join node. SortMergeJoin with an Exchange and Sort under both sides is the expensive default; BroadcastHashJoin with a BroadcastExchange under one side avoids the big shuffle; BroadcastNestedLoopJoin or CartesianProduct on two large inputs is the red flag. With AQE, check the final plan in the SQL tab of the Spark UI rather than the plan before execution.
Common mistakes
- Confusing join type with join strategy — A left join is not a strategy. It can run as a broadcast hash, shuffle hash, sort-merge or nested loop join.
- Forcing shuffle hash joins on skewed data — One hot key builds a hash table that does not fit, and the task fails instead of spilling.
- Not noticing a nested loop — A non-equi condition on two large tables quietly turns into rows × rows comparisons.
- Expecting a hint to always win — Spark drops hints that the join type or condition cannot support.
What you learned
- The difference between a join type and a join strategy
- The five physical join strategies and what each needs
- The order in which Spark picks a strategy
- How hints steer the choice, and how to fix a slow non-equi join
Key takeaways
- Join type is what you get; join strategy is how Spark computes it.
- Equi-joins can use broadcast hash, shuffle hash or sort-merge; non-equi joins only nested loops.
- Order: hints, broadcast hash, shuffle hash (if preferred), sort-merge, Cartesian, broadcast nested loop.
- Sort-merge is the default for large joins because it spills rather than failing.
- Rewrite OR and range conditions into equi-joins to escape nested loops.
Check yourself
4 questionsWhich strategy does a join on e.ts BETWEEN s.start AND s.end use when neither side is small?
Show the answer
CartesianProduct, for an inner join. There is no equality, so only nested-loop strategies apply; an inner join of two large sides becomes a Cartesian product.
Why does Spark prefer sort-merge over shuffle hash join by default?
Show the answer
It spills to disk under memory pressure, while a per-partition hash table must fit in memory. spark.sql.join.preferSortMergeJoin is true because sort-merge degrades gracefully.
Which join type can never use a broadcast hash join?
Show the answer
full outer. Both sides must keep unmatched rows, so neither can be the broadcast side.
Both sides of a join are hinted, one MERGE and one SHUFFLE_HASH. What wins?
Show the answer
MERGE. The priority is BROADCAST, MERGE, SHUFFLE_HASH, SHUFFLE_REPLICATE_NL.
Practice it
Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.
Keep going
Up next · lesson 8 of 30 · 3 min readBroadcast joins
Ship a small table to every executor and skip the shuffle, and when broadcasting backfires.
Related lessons
Previous: Narrow vs wide transformations
Primary sources: Join strategy hints · Performance tuning: join strategy hints