spark broadcast join threshold tuning
Tunes the Spark broadcast join threshold so small tables broadcast and large ones do not. Use when a join that should broadcast does not (slow sort-merge on a tiny dimension table), when a too-large broadcast OOMs the driver, or when you need per-query control of the threshold. Not for shuffle partition sizing, for AnalysisException errors, or for to_timestamp parse failures.
TL;DR
Raise spark.sql.autoBroadcastJoinThreshold when a small dimension table is not broadcasting, lower it when an oversized broadcast OOMs the driver, and use the broadcast() hint for per-query control. Broadcast is a driver-memory decision, not an executor one.
spark broadcast join threshold tuningUse this when
- A join on a small table runs as a slow sort-merge instead of broadcasting
- A broadcast join OOMs the driver
- You want one query to broadcast without changing the global setting
Not for this skill when
- Executors OOM during a shuffle, thats partition sizing
- Column resolution fails at analysis time
- Timestamps parse to null
Steps
- Check the current threshold and whether the small side actually fits:
print(spark.conf.get("spark.sql.autoBroadcastJoinThreshold"))
dim.count(), dim.rdd.map(lambda r: len(str(r))).sum() # rough size checkExpected output: the threshold in bytes (default 10MB, shown as 10485760) and a rough size of the dimension table. Broadcast is only considered when the small side is under the threshold.
- Raise the threshold when the dimension table is small but over 10MB:
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "104857600") # 100MBExpected output: set before the query. Joins against tables under 100MB now broadcast automatically. The broadcast data is collected on the driver first, so keep it comfortably under driver memory.
- Force or forbid broadcast per query with hints, without touching the global config:
from pyspark.sql.functions import broadcast
fact.join(broadcast(dim), "id") # force
fact.join(dim.hint("shuffle_hash"), "id") # forbid broadcast, use shuffle hashExpected output: the physical plan in the UI shows BroadcastHashJoin or SortMergeJoin as hinted. Hints are the right tool when one query is special.
- If the driver OOMs during broadcast, the table was bigger than you thought. Check the real size:
dim.cache().count()
# then check Storage tab in the UI for the cached sizeExpected output: the true in-memory size, often several times the on-disk size after deserialization. Lower the threshold or drop the hint so the join goes sort-merge.
- Confirm the plan in the UI or with explain:
fact.join(dim, "id").explain()Expected output: BroadcastHashJoin when broadcasting, SortMergeJoin otherwise. Trust the plan, not your assumption about which side Spark picked.
Variant phrasings
spark broadcast join not happening
Either the small side exceeds the threshold (step 2) or stats are missing so Spark overestimates its size. Run ANALYZE or check spark.sql.statistics handling.
spark driver out of memory broadcast
The broadcast table was too big for the driver (step 4). This is the failure mode of setting the threshold too high.
spark disable broadcast join
Set the threshold to -1 to disable auto-broadcast globally, or use the shuffle_hash / shuffle_merge hints per query.
Why it happens
A broadcast join sends the small table to every executor once, avoiding the shuffle of the big table entirely. Spark decides based on size estimates versus the threshold, collecting the small side on the driver. Wrong threshold in either direction gives you the worst of both: a shuffle you did not need, or a driver OOM.
Edge cases
- AQE can flip the join strategy at runtime based on actual sizes, keep it enabled and treat the threshold as a starting hint.
- Skewed big-side keys still hurt broadcast joins less than sort-merge, but the broadcast side must still fit in driver memory.
- Setting the threshold absurdly high (gigabytes) is a classic driver-OOM footgun, size it from measured table sizes.
spark.sql.broadcastTimeoutmatters when the broadcast collection is slow, not just large.
Provenance
Resolved from the public thread: https://vectle.com/posts/pstZoMrP9jMXELsxvIXUHVuQ
Maintainer review
No maintainer verification is recorded for this version.
This records the version a maintainer checked. It does not assert that the version is the latest upstream release.