Broadcast joins
Ship a small table to every executor and skip the shuffle, and when broadcasting backfires.
On this page
Show code in
Every code block on the page follows this.
You will learn
- How a broadcast hash join works
- When Spark broadcasts automatically, and how to force it
- Which join types can broadcast which side
- When broadcasting backfires
spark.sql.autoBroadcastJoinThreshold (10 MB), or when you add a broadcast hint.How it works
-
Collect
The small side is computed and collected to the driver
-
Broadcast
The driver sends it to every executor
-
Build
Each executor builds a hash table from it
-
Probe
Each task streams its large-side partition through the hash table
Compare a sort-merge join of two large tables: both sides shuffled by key, both sorted, then merged. Broadcasting replaces two shuffles and two sorts with one small transfer.
When Spark broadcasts
- Automatically, when its estimated size is below
spark.sql.autoBroadcastJoinThreshold(10 MB by default;-1disables it). - With AQE, at runtime, when the measured size after a shuffle is below
spark.sql.adaptive.autoBroadcastJoinThreshold, which defaults to the same value. - When you ask:
F.broadcast(df)or a SQL hint. Hints override the threshold.
from pyspark.sql import functions as F orders.join(F.broadcast(countries), "country_code")
SELECT /*+ BROADCAST(c) */ o.*, c.name FROM orders o JOIN countries c ON o.country_code = c.code
Which side can be broadcast
| Join type | Can broadcast |
|---|---|
| inner | Either side |
| left outer, left semi, left anti | The right side only |
| right outer | The left side only |
| full outer | Neither, as a broadcast hash join |
The rule: the side that must keep unmatched rows is the one streamed, so the other side is broadcast. Broadcasting the left side of a left join would lose track of left rows with no match across executors.
When broadcasting backfires
A "small" table that is not small
Too big for the hard limit
Timeouts
spark.sql.broadcastTimeout (300 seconds) fails the query, typically when the small side is itself an expensive query.Memory on every executor
Check what Spark chose
In explain(), look for BroadcastHashJoin with a BroadcastExchange under one side, versus SortMergeJoin with an Exchange on both sides. With AQE, the plan before execution may say SortMergeJoin and change at runtime; the SQL tab of the Spark UI shows the final plan.
Common mistakes
Broadcasting the preserved side of an outer join
Trusting size estimates after complex transformations
Setting the threshold to several GB
Key takeaways
- Broadcast joins ship the small side everywhere and avoid shuffling the large side.
- Automatic below 10 MB estimated, or at runtime with AQE; force with F.broadcast or a hint.
- Outer joins can only broadcast the non-preserved side; full outer cannot.
- Big broadcasts cost driver and executor memory and can time out.
Check yourself
3 questions1. What is the default spark.sql.autoBroadcastJoinThreshold?
Show the answer
10 MB. The default is 10 MB; -1 disables automatic broadcasting.
2. In a left outer join, which side can be broadcast?
Show the answer
Right. The left side must keep unmatched rows, so it is streamed and the right side is broadcast.
3. Which plan node shows a broadcast join?
Show the answer
BroadcastHashJoin. BroadcastHashJoin, with a BroadcastExchange under the small side.
Practice it
Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.
Go deeper
joinSpark internals
Adaptive Query ExecutionSpark internals
Dynamic partition pruning
Primary sources: Join strategy hints · Performance tuning