Apache Spark joins in depth: the five strategies and how one gets picked
Spark has five join strategies and one decision procedure that chooses between them. Here is how each works, the exact order Spark tries them in, which join types quietly forbid which strategy, and what to do per scenario.
- What is the difference between a join type and a join strategy?
- Architecture: what every join is made of
- Why are there five strategies, and not six?
- What are the five join strategies?
- How does Spark actually choose?
- How do you make each strategy appear?
- What does adaptive execution change about joins?
- How do you read the join out of a plan?
- Which join strategy fits which scenario?
- Common misconceptions
- Production tips
- Frequently asked questions
- Conclusion
- References
- Trademarks
TL;DR
- A join type is semantics (inner, outer, semi, anti). A join strategy is the algorithm (broadcast hash, shuffle hash, sort-merge, broadcast nested loop, cartesian). Mixing the two words up is why join tuning feels arbitrary.
- Spark’s choice is a documented, ordered procedure in
JoinSelection, not a cost model. Hints first, then broadcast, then shuffle hash, then sort-merge, then cartesian, then nested loop “as the final solution”.- The join type silently restricts the strategy. A full outer join can never be broadcast, and for a right outer join only the left side may be. Both are confirmed below by running them.
- Shuffle hash join needs three conditions at once, one of which is off by default, which is why you almost never see it without a hint.
- The 10 MB threshold is not the only size rule. Shuffle hash join has its own ceiling of
autoBroadcastJoinThresholdmultiplied byspark.sql.shuffle.partitions, which is 2 GB on the defaults, and a side effect: setting that threshold to-1disables shuffle hash join along with broadcasting.
Join tuning has a reputation for being folklore. Broadcast the small side, raise the threshold, add a hint, and if none of that works, salt the key. It works often enough to feel like knowledge and fails often enough to feel like luck.
The reason is that most of what circulates describes the five strategies without
describing the procedure that picks between them. That procedure is not
mysterious: it is written down, in order, in a comment above
JoinSelection in Spark’s own source, and once you have read it the behaviour
stops being surprising. A full outer join that refuses to broadcast is not a
bug, and no amount of raising the threshold will change it.
This post has three parts. What each strategy actually does, the exact order Spark tries them in, and which scenario should get which.
What is the difference between a join type and a join strategy?
This is the distinction to get right first, because the two are chosen independently and the first one constrains the second.
| Join type | Join strategy | |
|---|---|---|
| Decides | Which rows appear in the result | How the rows are brought together |
| You write it | Yes, it is your query’s meaning | No, Catalyst picks it |
| Examples | inner, left outer, full outer, left semi, left anti, cross | broadcast hash, shuffle hash, sort-merge, broadcast nested loop, cartesian |
| Changes results | Yes | Never |
A strategy change can make a query ten times faster or fail outright, and it cannot change the answer. A type change always changes the answer. So when someone says “we changed the join and it got faster”, the useful question is which of the two they changed.
The types, briefly, since the rest of the post assumes them:
trips.join(cities, "city_id") # inner, the default
trips.join(cities, "city_id", "left") # left outer
trips.join(cities, "city_id", "full") # full outer
trips.join(cities, "city_id", "left_semi") # left rows that match
trips.join(cities, "city_id", "left_anti") # left rows that do not
trips.crossJoin(cities) # every pair
Semi and anti deserve a word because they are under-used. Both filter the left
side by whether a match exists and return no columns from the right, so a
left row matching a hundred right rows is emitted once. An inner join in the
same position emits it a hundred times and needs a distinct afterwards, which
is a second shuffle.
Architecture: what every join is made of
All five strategies are assembled from the same three parts, and naming them makes the differences between the strategies easy to hold.
flowchart TB
subgraph P["every join has these"]
A["<b>build side</b><br/>materialised into a lookup<br/>structure, or sorted"]
B["<b>stream side</b><br/>read once and probed<br/>against the build side"]
C["<b>co-location step</b><br/>how matching keys are made<br/>to meet on one executor"]
end
A --> J["the join operator"]
B --> J
C --> J
The build side is the one Spark materialises: a hash table for the hash joins, a sorted run for sort-merge. It is the side that costs memory, which is why “which side is the build side” is the question behind most join failures.
The stream side is read once and probed. It costs nothing to hold, which is why you want the large table here.
The co-location step is where the strategies genuinely differ, and it is the expensive part:
| Strategy | How keys are co-located | Network cost |
|---|---|---|
| Broadcast hash | Copy the whole build side everywhere | One copy per executor |
| Shuffle hash | Shuffle both sides by key | Full shuffle |
| Sort-merge | Shuffle both sides by key, then sort | Full shuffle plus sort |
| Broadcast nested loop | Copy one side everywhere, compare all pairs | One copy per executor |
| Cartesian | Replicate partitions across each other | Every partition pair |
Read down that column and the whole subject compresses to one trade: either you copy the small side to where the data already is, or you move both sides to a common place. Broadcasting does the first, the shuffle-based strategies do the second, and bucketing is how you pay for the second once at write time instead of on every read.
Why are there five strategies, and not six?
The five strategies are not five unrelated algorithms. They are combinations of two independent choices, and separating them is the clearest way to hold the whole subject.
Choice one, the data exchange: how matching keys are made to meet on one executor.
- Broadcast. Copy one side to every executor. The other side never moves.
- Shuffle. Repartition both sides by the join key so matching keys land in the same partition.
Choice two, the join algorithm: what happens once the rows are together.
- Hash. Build a hash table from one side, probe it with the other.
- Sort-merge. Sort both sides, walk them with a cursor each.
- Nested loop. Compare every pair, evaluating an arbitrary condition.
Two exchanges times three algorithms is six combinations. Spark ships five physical operators, and the grid shows exactly which cell is missing:
| Hash | Sort-merge | Nested loop | |
|---|---|---|---|
| Broadcast | BroadcastHashJoinExec |
does not exist | BroadcastNestedLoopJoinExec |
| Shuffle | ShuffledHashJoinExec |
SortMergeJoinExec |
CartesianProductExec |
Why there is no broadcast sort-merge join is worth a moment, because it explains what sorting is actually for. Sorting buys bounded memory: you never hold a whole side, and the sort can spill. But if a side has already been copied to every executor, it is in memory by definition, so a hash table over it costs nothing extra and gives constant-time probes. Sorting both sides to get the same answer more slowly has no upside. Sorting only earns its place when you had to shuffle anyway and cannot afford to hold a side in memory.
So the decision Spark makes is really two decisions:
flowchart TB
A["can one side be copied<br/>to every executor?"]
A -->|"yes"| B{"equi-join?"}
A -->|"no"| C{"equi-join?"}
B -->|"yes"| D["BroadcastHashJoin"]
B -->|"no"| E["BroadcastNestedLoopJoin"]
C -->|"yes"| F{"hold a side in memory<br/>per partition?"}
C -->|"no"| G["CartesianProduct<br/><i>inner-like only</i>"]
F -->|"yes"| H["ShuffledHashJoin"]
F -->|"no"| I["SortMergeJoin"]
Read it as: the exchange is decided by size, and the algorithm by the condition and by how much memory you can hold. An equi-join can hash or sort-merge; a non-equi join has only nested loop, which is why it collapses to two options. That single sentence predicts most of what the planner does.
What are the five join strategies?
Spark’s source documents all five with their limits. The table below is that comment, condensed:
| Strategy | Equi-join only | Keys must be sortable | Join types |
|---|---|---|---|
| Broadcast hash join | Yes | No | All except full outer |
| Shuffle hash join | Yes | No | All |
| Sort-merge join | Yes | Yes | All |
| Broadcast nested loop | No | No | All, but efficient only in some combinations |
| Cartesian product | No | No | Inner-like only |
An equi-join is one whose condition is equality on keys, a.id = b.id.
Anything else, a range or an inequality, is a non-equi join, and only the bottom
two rows can execute it at all. That single fact explains most surprises in the
plan.
Broadcast hash join
Spark collects the small side to the driver, sends a copy to every executor, and each executor builds a hash table from it. Each partition of the large side then probes that table locally. The large side is never shuffled.
flowchart TB
S["small side: cities"] --> D["driver collects it"]
D --> BC[("broadcast:<br/>one copy per executor")]
BC --> E1["executor 1<br/>hash table"]
BC --> E2["executor 2<br/>hash table"]
BC --> E3["executor 3<br/>hash table"]
T1["trips partition 0<br/>does not move"] --> E1
T2["trips partition 1<br/>does not move"] --> E2
T3["trips partition 2<br/>does not move"] --> E3
E1 --> R["output, no shuffle<br/>of the large side"]
E2 --> R
E3 --> R
Cost: one copy of the small side per executor, held in memory for the duration, plus the driver having to assemble it first.
That driver step is the part people forget. A broadcast that is too large fails
on the driver, not the executors, and the error often arrives as a
spark.sql.broadcastTimeout after the default 300 seconds rather than as an
obvious out-of-memory.
Key points
- Build side: the small one, if
canBroadcastBySizepasses and the join type allows that side. - Network: one copy of the small side per executor. The large side never moves.
- Memory: the hash table sits on every executor, and on the driver first.
- Wins when: one side is genuinely small after filters and projections are applied.
- Fails when: the small side is not small, or the join type forbids that build side.
- In the plan:
BroadcastHashJoinwith aBroadcastExchangechild.
Shuffle hash join
Both sides are shuffled on the join key. Then, for each partition, Spark builds a hash table from the smaller side and probes it with the larger. No sorting.
flowchart TB
A["riders<br/>the smaller side"] --> SH[("shuffle both sides<br/>by the join key")]
B["trips<br/>the larger side"] --> SH
SH --> P0["partition 0<br/>build hash from riders<br/>probe with trips"]
SH --> P1["partition 1<br/>build hash from riders<br/>probe with trips"]
SH --> P2["partition 2<br/>build hash from riders<br/>probe with trips"]
P0 --> O["output, never sorted"]
P1 --> O
P2 --> O
It sits in a narrow band: one side too big to broadcast, but small enough that a per-partition hash table fits. It skips the sort that sort-merge pays, and in exchange it can run out of memory where sort-merge would spill.
Key points
- Build side: the smaller side, per partition, not the whole table.
- Network: a full shuffle of both sides, the same as sort-merge.
- Memory: one hash table per partition, which is what can blow up.
- Wins when: the build side is much smaller but still too big to broadcast, and the sort is the thing you want to skip.
- Fails when: a partition’s build side does not fit, since a hash table cannot spill the way a sort can.
- In the plan:
ShuffledHashJoinwith anExchange hashpartitioningunder each side and noSortnode.
Sort-merge join
Both sides are shuffled on the key, each partition is sorted, and the two sorted streams are walked together with a cursor on each.
flowchart TB
A["trips"] --> S[("shuffle both sides<br/>by the join key")]
B["cities"] --> S
S --> S0["partition 0<br/>sort both sides"]
S --> S1["partition 1<br/>sort both sides"]
S0 --> M0["merge with one<br/>cursor per side"]
S1 --> M1["merge with one<br/>cursor per side"]
M0 --> O["output, sorted by key"]
M1 --> O
The merge itself is the cheap part. Each side has a cursor, and whichever cursor points at the smaller key advances:
left (sorted): a a b c d ...
right (sorted): a b b d e ...
^ ^
Its virtue is bounded memory. Only the merge buffers are held, not a whole hash
table, and the sort can spill, so a partition larger than memory still finishes.
That is why it is the default for large-to-large joins, and why
spark.sql.join.preferSortMergeJoin defaults to true.
The cost is the sort, and the requirement is that the keys are sortable. It is the only strategy with that constraint.
Key points
- Build side: neither, in the hash-table sense. Both sides are sorted and streamed.
- Network: a full shuffle of both sides, plus a sort per partition.
- Memory: bounded, and the sort spills to disk rather than failing.
- Wins when: both sides are large, which is why it is the default.
- Fails when: rarely, and not on memory. Its failure mode is slowness from the shuffle and the sort, or one skewed partition, rather than an out-of-memory.
- In the plan:
SortMergeJoinwith aSortnode above anExchangeon each side.
Broadcast nested loop join
One side is broadcast, and then every row of one side is compared against every row of the other, evaluating an arbitrary condition. It is the fallback that makes non-equi joins possible at all.
The source notes it is optimised for three specific shapes, and slow otherwise:
- Broadcasting the left side in a right outer join.
- Broadcasting the right side in a left outer, left semi, left anti or existence join.
- Broadcasting either side in an inner-like join.
flowchart TB
S["one side, broadcast whole"] --> E1["executor 1"]
S --> E2["executor 2"]
P1["partition of the other side"] --> E1
P2["partition of the other side"] --> E2
E1 --> C1["compare every local row<br/>against every broadcast row"]
E2 --> C2["compare every local row<br/>against every broadcast row"]
C1 --> O["rows where the<br/>condition holds"]
C2 --> O
Outside those, “we need to scan the data multiple times, which can be rather
slow”. Seeing BroadcastNestedLoopJoin on two large inputs means the query is
doing something close to a cross product.
Key points
- Build side: the broadcast one, and only the shapes listed above are efficient.
- Network: one copy of the broadcast side per executor.
- Cost: the row comparisons, which grow as the product of the two row counts.
- Wins when: the condition is not an equality and one side is genuinely tiny.
- Fails when: both sides are large, where it is the “may OOM” last resort.
- In the plan:
BroadcastNestedLoopJoinwith the condition printed inline.
Cartesian product
Also called shuffle-and-replicate nested loop. Every partition of one side is
paired with every partition of the other, producing the full cross product. It
supports inner-like joins only, and it is what you get from a deliberate
crossJoin once broadcasting is off the table.
flowchart TB
A0["A partition 0"] --> X["every A partition paired<br/>with every B partition"]
A1["A partition 1"] --> X
B0["B partition 0"] --> X
B1["B partition 1"] --> X
X --> P00["A0 x B0"]
X --> P01["A0 x B1"]
X --> P10["A1 x B0"]
X --> P11["A1 x B1"]
With two partitions a side that is four tasks. With 200 a side it is 40,000, and that multiplication is the whole reason this strategy is a warning sign.
Key points
- Build side: neither. Partitions are replicated against each other.
- Network: every partition pair, so the shuffle grows as the product of the partition counts.
- Cost: the output itself, which is the product of the two row counts.
- Wins when: you actually want a cross product, on small inputs.
- Fails when: you did not want one, which is the usual case.
- In the plan:
CartesianProduct, almost always worth investigating.
How does Spark actually choose?
Here is the part that turns folklore into a procedure. JoinSelection follows a
fixed order, and the first applicable rule wins.
For an equi-join, hints are consulted first, in this order:
BROADCAST: pick broadcast hash join if the join type supports it.MERGE: pick sort-merge join if the keys are sortable.SHUFFLE_HASH: pick shuffle hash join if the join type supports it.SHUFFLE_REPLICATE_NL: pick cartesian product if the join type is inner-like.
With no hints, for an equi-join:
- Broadcast hash join, if one side is small enough to broadcast and the join type allows that side to be the build side.
- Shuffle hash join, if one side can build a local hash map, is much smaller than the other, and
spark.sql.join.preferSortMergeJoinisfalse. - Sort-merge join, if the join keys are sortable.
- Cartesian product, if the join type is inner-like.
- Broadcast nested loop join as the last resort. The source is blunt about it: “It may OOM but we don’t have other choice.”
For a non-equi join the list is much shorter, which is why these are so often slow:
- Broadcast nested loop join, if one side is small enough to broadcast.
- Cartesian product, if the join type is inner-like.
- Broadcast nested loop join anyway.
flowchart TB
H{"a join hint?"} -->|"yes"| HH["honour it if the type<br/>and keys allow"]
H -->|"no"| EQ{"equi-join?"}
EQ -->|"no"| NE["broadcast nested loop if a side is small,<br/>else cartesian if inner-like,<br/>else nested loop anyway"]
EQ -->|"yes"| B{"a side small enough,<br/>and allowed as build side?"}
B -->|"yes"| BHJ["BroadcastHashJoin"]
B -->|"no"| SH{"local hash map fits,<br/>much smaller,<br/>preferSortMergeJoin false?"}
SH -->|"yes"| SHJ["ShuffledHashJoin"]
SH -->|"no"| SM{"keys sortable?"}
SM -->|"yes"| SMJ["SortMergeJoin"]
SM -->|"no"| C["cartesian if inner-like,<br/>else nested loop"]
The size rules, exactly
Three predicates decide “small enough”, and only the first is widely known.
Broadcast. canBroadcastBySize is sizeInBytes >= 0 && sizeInBytes <=
autoBroadcastJoinThreshold, where the threshold is
spark.sql.autoBroadcastJoinThreshold, default 10485760, which is 10 MB.
There is a detail worth having: when the statistics are runtime statistics,
meaning adaptive execution has measured a completed stage, Spark uses
spark.sql.adaptive.autoBroadcastJoinThreshold instead, falling back to the
static one when that is unset. So plan-time and runtime broadcasting can be
governed by two different numbers.
Shuffle hash, ceiling. canBuildLocalHashMapBySize is
sizeInBytes < autoBroadcastJoinThreshold * numShufflePartitions. On the
defaults that is 10 MB multiplied by 200:
10,485,760 x 200 = 2,097,152,000 bytes = ~2 GB
So the shuffle-hash ceiling is roughly 2 GB by default, two hundred times the broadcast threshold, and it moves when you change the shuffle partition count.
Shuffle hash, ratio. muchSmaller is
a.size * spark.sql.shuffledHashJoinFactor <= b.size, and that factor defaults
to 3. The source explains why: “The cost to build hash map is higher than
sorting, we should only build hash map on a table that is much smaller than
other one.”
So shuffle hash join requires three things simultaneously: under the 2 GB
ceiling, at least three times smaller than the other side, and
preferSortMergeJoin flipped to false. That conjunction, with the last
condition off by default, is why it is rare in practice.
Which build side each join type allows
This is the rule that makes broadcast joins look unpredictable, and it is pure lookup, not heuristics.
| Join type | Broadcast the left? | Broadcast the right? |
|---|---|---|
| Inner, cross | Yes | Yes |
| Left outer | No | Yes |
| Right outer | Yes | No |
| Full outer | No | No |
| Left semi, left anti | No | Yes |
The logic follows from correctness. A left outer join must emit every left row, so the left side has to be streamed and the right one built. A right outer is the mirror image. A full outer must emit unmatched rows from both sides, so neither can be the build side, and no full outer join is ever a broadcast hash join regardless of how small either side is.
Shuffle hash join is less restricted, which is why the strategy table says it supports all join types: its build side may be the left for inner, left outer, full outer and right outer, and the right for those plus left semi, left anti and existence joins.
How do you make each strategy appear?
Every rule above is testable. This section builds two tables, forces all five
strategies, and reads back the one that actually executed. The transcripts below
are from Spark 4.1.3 with spark.sql.shuffle.partitions set to 8, and the
five selection predicates quoted earlier were checked against the 4.1.3 source.
A fact table and a dimension:
from pyspark.sql import SparkSession, functions as F
spark = (SparkSession.builder
.appName("join-strategies")
.master("local[4]")
.config("spark.sql.shuffle.partitions", "8")
.getOrCreate())
cities = (spark.range(0, 600)
.select(F.col("id").alias("city_id"),
F.concat(F.lit("city-"), F.col("id")).alias("city_name"),
F.when(F.col("id") % 3 == 0, "IN")
.when(F.col("id") % 3 == 1, "US")
.otherwise("DE").alias("country_code")))
trips = (spark.range(0, 3_000_000)
.select(F.col("id").alias("trip_id"),
(F.col("id") % 500).alias("city_id"),
(F.col("id") % 90000).alias("rider_id"),
(F.col("id") % 4000 / 100.0).cast("decimal(10,2)").alias("fare_amount")))
Two helpers. The first reports the strategy that ran, the second reports the size Catalyst believes a side to be, which is the input to every size rule:
def strategy(df):
plan = df._jdf.queryExecution().executedPlan().toString()
for name in ["BroadcastHashJoin", "ShuffledHashJoin", "SortMergeJoin",
"BroadcastNestedLoopJoin", "CartesianProduct"]:
if name in plan:
return name
return "unknown"
def size_in_bytes(df):
return int(df._jdf.queryExecution().optimizedPlan().stats().sizeInBytes())
For these two tables the estimates are:
size_in_bytes(cities) = 16,800 bytes
size_in_bytes(trips) = 60,000,000 bytes
trips references only the first 500 city_id values, so 100 of the 600 cities
have no trips. That matters for the outer joins later. cities is far under the
10 MB broadcast threshold and trips is far over it, so the defaults should
broadcast. Now force each strategy in turn:
# 1 broadcast hash join, the default here because cities is under the threshold
strategy(trips.join(cities, "city_id"))
# 2 sort-merge, by taking broadcasting away
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")
strategy(trips.join(cities, "city_id"))
spark.conf.unset("spark.sql.autoBroadcastJoinThreshold")
# 3 shuffle hash, by asking for it
strategy(trips.join(cities.hint("SHUFFLE_HASH"), "city_id"))
# 4 broadcast nested loop, by using an inequality instead of an equality
strategy(trips.join(cities, trips.city_id > cities.city_id))
# 5 cartesian, a deliberate cross product with broadcasting off
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")
strategy(trips.limit(1000).crossJoin(cities.limit(50)))
spark.conf.unset("spark.sql.autoBroadcastJoinThreshold")
The operator line each one produced:
1 BroadcastHashJoin [city_id#6L], [city_id#1L], Inner, BuildRight, false
2 SortMergeJoin [city_id#6L], [city_id#1L], Inner
3 ShuffledHashJoin [city_id#6L], [city_id#1L], Inner, BuildRight
4 BroadcastNestedLoopJoin BuildRight, Inner, (city_id#6L > city_id#1L)
5 CartesianProduct
All five, from one dataset, by changing one thing at a time. Note that case 4 returned 748,500,000 rows from a 3,000,000-row table: an inequality against 600 cities is a near-cross-product, which is what the strategy name is telling you.
Each of the four hints behaves as the procedure says, overriding the size rules but not the join-type rules:
| Hint | Written as | Strategy it produced |
|---|---|---|
BROADCAST |
cities.hint("BROADCAST") |
BroadcastHashJoin |
MERGE |
cities.hint("MERGE") |
SortMergeJoin |
SHUFFLE_HASH |
cities.hint("SHUFFLE_HASH") |
ShuffledHashJoin |
SHUFFLE_REPLICATE_NL |
cities.hint("SHUFFLE_REPLICATE_NL") |
CartesianProduct |
A hint is the one lever that skips the size checks entirely, which is why
BROADCAST on a large table moves the failure to the driver rather than
preventing it. The last row is the most destructive of the four: it printed
CartesianProduct (city_id#6L = city_id#1L), meaning the equality was demoted
from a join key to a filter applied after the cross product was built.
The join type overrides the size
The size rules are only consulted after the join type has had its say, and this is where the surprises live. One rule sits underneath all of them: an outer join must emit every row of its preserved side, so that side has to be streamed, and the build side is always the other one.
A left outer preserves the left and therefore builds the right. A right outer preserves the right and builds the left. A full outer preserves both, so it can build neither. That is the entire build-side table, derived rather than memorised.
Four joins over the same two tables make the consequence concrete:
strategy(trips.join(cities, "city_id", "right")) # preserves cities, the small side
strategy(cities.join(trips, "city_id", "left")) # preserves cities as well
strategy(cities.join(trips, "city_id", "right")) # preserves trips, the large side
strategy(trips.join(cities, "city_id", "left")) # preserves trips as well
trips RIGHT JOIN cities 3,000,100 rows SortMergeJoin [city_id#6L], [city_id#1L], RightOuter
cities LEFT JOIN trips 3,000,100 rows SortMergeJoin [city_id#1L], [city_id#6L], LeftOuter
cities RIGHT JOIN trips 3,000,000 rows BroadcastHashJoin [city_id#1L], [city_id#6L], RightOuter, BuildLeft, false
trips LEFT JOIN cities 3,000,000 rows BroadcastHashJoin [city_id#6L], [city_id#1L], LeftOuter, BuildRight, false
The first two lines are one query written two ways. Both preserve cities, both
return 3,000,100 rows, and comparing them row by row finds no difference in
either direction. Both sort-merge. The last two are also one query written two
ways, both preserve trips, and both broadcast.
The two pairs, though, are not the same query. They differ by exactly the 100 cities that have no trips. That is the trap in the advice usually given here: swapping the argument order of an outer join changes which side is preserved, so it changes the answer. It is a join-type change wearing the costume of a refactor, and it is the one rewrite that can quietly drop rows.
So the honest rule is narrower than “swap the sides”:
- If you must preserve the large side, the build side is the small one, and the broadcast comes for free.
- If you must preserve the small side, the build side is the large one, and no threshold, hint or reordering produces a broadcast hash join. Sort-merge is the correct plan, and the lever you have is reducing the large side, not persuading the optimiser.
Two more results from the same tables close the picture:
trips FULL JOIN cities 3,000,100 rows SortMergeJoin ... FullOuter
trips LEFT ANTI JOIN cities 0 rows BroadcastHashJoin ... LeftAnti, BuildRight, false
A full outer preserves both sides, so it sort-merges however small cities is.
A left anti preserves the left and returns no columns from the right at all, so
the right is always available as the build side and it broadcasts. It returns
zero rows here because every trip does have a matching city, which is the
correct answer to “trips whose city is missing”.
Why shuffle hash join is rare, demonstrated
The three conditions are easier to believe when you watch them fail one at a
time. Take a 400,000-row riders table against trips on rider_id:
riders = (spark.range(0, 400_000)
.select(F.col("id").alias("rider_id"),
F.concat(F.lit("rider-"), F.col("id")).alias("rider_name")))
def attempt(threshold, prefer_smj):
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", str(threshold))
spark.conf.set("spark.sql.join.preferSortMergeJoin", prefer_smj)
return strategy(trips.select("trip_id", "rider_id").join(riders, "rider_id"))
attempt(-1, "false") # the usual "disable broadcast" trick
attempt(2 * 1024**2, "false") # broadcast still off, but a positive threshold
attempt(2 * 1024**2, "true") # preferSortMergeJoin back at its default
With riders at 7,200,000 bytes and the projected trips at 36,000,000:
| Attempt | canBroadcastBySize |
Hash-map ceiling | canBuildLocalHashMapBySize |
muchSmaller |
preferSortMergeJoin |
Strategy |
|---|---|---|---|---|---|---|
threshold = -1 |
false | -1 x 8 = -8 |
false | true | false | SortMergeJoin |
threshold = 2 MB |
false | 2 MB x 8 = 16 MB |
true | true | false | ShuffledHashJoin |
threshold = 2 MB, prefer on |
false | 2 MB x 8 = 16 MB |
true | true | true | SortMergeJoin |
The middle row is a shuffle hash join with no hint at all, which is the thing most people have never seen. The other two rows each break one condition.
The first row deserves its own warning, because it is a trap in the most common
tuning gesture there is. Setting autoBroadcastJoinThreshold to -1 does not
just disable broadcasting, it also makes shuffle hash join unreachable, since
canBuildLocalHashMapBySize multiplies that same threshold by the shuffle
partition count and -1 x 8 is negative. Every size below zero fails the test.
If you want to stop broadcasting but keep shuffle hash join available, set the
threshold to a small positive number instead of -1.
What does adaptive execution change about joins?
Everything above is plan-time. Adaptive query execution re-plans after a shuffle
stage completes, when the sizes are measurements rather than estimates, and
spark.sql.adaptive.enabled defaults to true.
For joins it does two things:
- Demotes a sort-merge join to a broadcast hash join when a side turns out small enough, using
spark.sql.adaptive.autoBroadcastJoinThreshold. - Splits skewed partitions, with
spark.sql.adaptive.skewJoin.enableddefaulttrue.
The skew thresholds are strict and worth knowing, because they explain “skew
handling is on but nothing happened”. A partition is treated as skewed only when
it is larger than spark.sql.adaptive.skewJoin.skewedPartitionFactor, default
5.0, times the median and larger than
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes, default 256MB.
On a job whose partitions are all under 256 MB, it never fires no matter how
uneven they are.
The practical consequence for reading plans: explain() shows the planned
shape, not the final one. Before execution the tree is wrapped in
AdaptiveSparkPlan with isFinalPlan=false. The SQL tab of the UI after the
query finishes shows what actually ran.
How do you read the join out of a plan?
df.explain("formatted")
Here is the real output for the broadcast join above, projected down to three columns. The tree comes first, then a numbered detail block per operator:
== Physical Plan ==
AdaptiveSparkPlan (9)
+- Project (8)
+- BroadcastHashJoin Inner BuildRight (7)
:- Project (3)
: +- Filter (2)
: +- Range (1)
+- BroadcastExchange (6)
+- Project (5)
+- Range (4)
(6) BroadcastExchange
Input [2]: [city_id#1L, city_name#2]
Arguments: HashedRelationBroadcastMode(List(input[0, bigint, false]),false), [plan_id=28]
(7) BroadcastHashJoin
Left keys [1]: [city_id#6L]
Right keys [1]: [city_id#1L]
Join type: Inner
Join condition: None
(9) AdaptiveSparkPlan
Output [3]: [trip_id#5L, city_name#2, fare_amount#8]
Arguments: isFinalPlan=false
Four things in that output are worth reading deliberately. BuildRight on node
7 names the build side. BroadcastExchange on node 6 appears under the right
branch only, so the left branch has no exchange and the large side is not
shuffled. Node 6 lists city_id and city_name but not country_code, because
the projection pushed down and shrank what gets broadcast. And node 9 says
isFinalPlan=false, so this is the plan before adaptive execution has had a say.
What to look for, in order of how much it tells you:
| In the plan | Meaning |
|---|---|
BroadcastHashJoin |
The small side is being broadcast. Usually what you want |
SortMergeJoin |
Both sides shuffled and sorted. Correct for large-to-large |
ShuffledHashJoin |
Rare without a hint |
BroadcastNestedLoopJoin |
Non-equi, or a missing condition |
CartesianProduct |
Almost always a mistake |
Exchange hashpartitioning(k, 200) |
The shuffle, with its key and partition count |
BroadcastExchange |
The broadcast being built |
ReusedExchange |
One shuffle serving two branches. Good news |
isFinalPlan=false |
AQE has not re-planned yet |
A useful habit: count the Exchange nodes. That is the number of shuffles, and
if two consecutive operations key on the same column, one exchange can serve
both.
Which join strategy fits which scenario?
| Scenario | Aim for | How |
|---|---|---|
| Large fact, small dimension | Broadcast hash | Nothing, if the dimension is under 10 MB. Filter and project it first if not |
| Large fact, medium dimension (tens of MB) | Broadcast hash | Reduce it first, then raise the threshold deliberately if it is genuinely small in memory |
| Large to large, recurring | Sort-merge, or no shuffle at all | Bucket both tables on the key with the same bucket count |
| Large to large, one-off | Sort-merge | Filter and project both sides, size the shuffle, handle skew |
| Full outer with a small side | Sort-merge, unavoidable | Accept it, or restructure as two outer joins and a union |
| Outer join preserving the large side | Broadcast hash | Nothing. The build side is the small one already |
| Outer join preserving the small side | Sort-merge, unavoidable | The preserved side cannot be built. Reduce the other side instead |
| Range or inequality condition | Broadcast nested loop | Add an equality that buckets both sides, so the nested loop runs within buckets |
| Existence check | Semi or anti join | left_semi or left_anti, never inner plus distinct |
| Skewed join key | Sort-merge plus skew handling | AQE first, then salt the hot key |
Two of those need expanding.
The range join. A condition like ON t.ts BETWEEN d.start AND d.end has no
equality, so Spark can only nested-loop it. The standard fix is to manufacture a
key: bucket both sides by a coarser value, join on that equality and the
range, so the comparison happens within each bucket rather than across the whole
product.
(events.withColumn("day", F.to_date("ts"))
.join(windows.withColumn("day", F.to_date("start")),
["day"]) # an equality to partition on
.filter((F.col("ts") >= F.col("start")) & (F.col("ts") <= F.col("end"))))
The full outer with a small side. If broadcasting matters more than the single scan, a full outer can be rewritten as a left outer union an anti join the other way, both of which can broadcast. It is more code and two passes, so only worth it when the sort-merge is genuinely the bottleneck.
Common misconceptions
“Raising autoBroadcastJoinThreshold will make this broadcast.” Not if the
join type forbids it. A full outer join can never be a broadcast hash join,
because both sides must emit unmatched rows and neither may be the build side.
No threshold changes that. Check the join type before the size.
“Swapping the sides of an outer join is a free optimisation.” It changes the answer. An outer join preserves one side, and swapping the argument order swaps which side that is. Measured on 600 cities where 100 have no trips, the two orderings differ by exactly those 100 rows, while returning the same count on data where everything matches, which is how the bug hides.
“Setting autoBroadcastJoinThreshold to -1 only turns off broadcasting.”
It also removes shuffle hash join from consideration. The shuffle-hash ceiling
is autoBroadcastJoinThreshold * spark.sql.shuffle.partitions, so -1 makes it
negative and every candidate fails the size test. Use a small positive value if
you want to stop broadcasting and keep the option.
“A different join strategy can change my results.” It cannot. Strategies are algorithms for the same semantics. If the rows changed, the join type or the condition changed, or nulls in the key are behaving as SQL says they should.
Production tips
- Read the plan before tuning. The strategy Spark chose tells you which rule fired, and therefore which lever exists.
- Reduce before joining. Filtering and projecting a side is the highest-value change available, because it can move the join into the broadcast band and change the strategy entirely.
- Do not hint
BROADCASTon something large. A hint overrides the size check, so you move the failure to the driver, which must collect it first. - Prefer a small positive
autoBroadcastJoinThresholdover-1. Setting it to-1also drives the shuffle-hash ceiling negative, removing that strategy from consideration along with broadcasting. - Treat a full outer join as unbroadcastable. Raising the threshold will never help; the join type is what forbids it.
- For an outer join, ask which side is preserved, not which side is small. The build side is always the non-preserved one, so an outer join that must preserve the small side cannot broadcast at any threshold. Reordering the arguments changes the result, not just the plan.
- Fix the estimate rather than forcing the plan where you can, with
ANALYZE TABLE ... COMPUTE STATISTICS, so Catalyst chooses correctly on its own. - Bucket the keys you join on repeatedly. It is the only way to remove the shuffle from a large-to-large join permanently.
- Trust the SQL tab over
explain()for what actually ran, since AQE re-plans after the shuffle.
Frequently asked questions
Why is my small table not being broadcast? Four common reasons: the estimate is of the unfiltered table because the filter did not push down; the table is compressed on disk and larger in memory than you think; there are no statistics so Spark fell back to file size; or the join type forbids that side as the build side. Check the last one first, because no config change fixes it.
Does a bigger autoBroadcastJoinThreshold make joins faster?
Sometimes, and it moves the failure mode from slow to fragile. Every executor
holds a full copy and the driver assembles it first, so a 500 MB broadcast across
50 executors is 25 GB of cluster memory plus a driver that must hold it. Raise it
deliberately, with a number you can justify.
Why do I almost never see ShuffledHashJoin?
Because it needs three conditions at once and one of them,
spark.sql.join.preferSortMergeJoin being false, is not the default. Spark
prefers sort-merge because it cannot run out of memory, while a hash join can.
There is also a trap: if you disabled broadcasting with
autoBroadcastJoinThreshold = -1, the shuffle-hash ceiling became negative too,
so the strategy cannot be chosen without a hint no matter what else you set. The
table above shows the same join picking it once the threshold is a small
positive number.
Can a join strategy change my results? No. Strategies are algorithms for the same semantics. If results changed, the join type or the condition changed, or there are nulls in the key behaving as SQL says they should.
Why is my join a CartesianProduct when I wrote a condition?
Because Catalyst could not use the condition as an equality: it is an
inequality, it is wrapped in an expression it cannot match on, or it sits in a
where after a cross join in a way that could not be pushed into the join.
What about joining on a collated string column?
In Spark 4 this matters. Hash joins require keys that are binary-stable, and a
case- or accent-insensitive collation is not, because two rows can be equal
without having equal bytes. Spark does not give up on hashing for that, though:
it rewrites each key into collationkey(k), a binary-stable form, and hashes
that. On Spark 4.1.3 a join on two UTF8_LCASE columns broadcast by default,
took ShuffledHashJoin under a SHUFFLE_HASH hint, and sort-merged only once
broadcasting was off, with both join keys printed as collationkey(k) in
every plan. Hash joins stay available; the price is a collationkey
evaluated on every row of both sides.
Conclusion
The thing worth taking away is that join selection is a lookup, not a judgement. Spark asks, in a fixed order: is there a hint, is this an equi-join, is one side small enough, does the join type permit that side as the build side, are the keys sortable. The first rule that fits wins. Nothing about it is adaptive except the parts AQE re-plans after a shuffle.
That reframes tuning. You are not persuading an optimiser, you are changing the inputs to a decision procedure. Filtering a dimension until it fits under 10 MB changes the answer to “is one side small enough”. Deciding which side an outer join preserves settles which side may be built, though it changes the result too, so it is a modelling decision rather than a tuning knob. Bucketing changes whether a shuffle is needed at all. Each of those is a specific rule with a specific lever.
And three answers are simply no. A full outer join will not broadcast, a non-equi join will not hash, and an outer join that must preserve its small side will not broadcast either. When you hit those, the fix is to change the query shape, or to accept the plan and shrink the other side, rather than to keep raising thresholds.
References
SparkStrategies.scala, whoseJoinSelectioncomment documents the five strategies and the order they are tried injoins.scalafor the size predicates and the build-side rules per join type- SQL performance tuning for join hints and adaptive query execution
- Apache Spark architecture for the shuffle and stage machinery these strategies sit on
- Two community treatments this post builds on: Different types of Spark join strategies for separating data exchange from join algorithm, and Join strategies in Apache Spark, a hands-on approach for forcing each strategy on one dataset
Trademarks
Apache Spark, Apache Hive, Apache Parquet and Apache are either registered trademarks or trademarks of The Apache Software Foundation in the United States and other countries.
Found this useful?
These posts and tools are free. If one saved you an afternoon, you can buy me a coffee.