Skip to content
Good engineers know 3 min read · spark.sql.autoBroadcastJoinThreshold 1 practice problem ↓

Broadcast joins

Ship a small table to every executor and skip the shuffle, and when broadcasting backfires.

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

Read first

Comfortable with these? Read on.

TL;DR If one side of a join is small, Spark sends a full copy to every executor and joins the large side in place, with no shuffle. It happens automatically below spark.sql.autoBroadcastJoinThreshold (10 MB), or when you add a broadcast hint.

How it works

  1. Collect

    The small side is computed and collected to the driver

  2. Broadcast

    The driver sends it to every executor

  3. Build

    Each executor builds a hash table from it

  4. 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; -1 disables 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.
PySparkSpark SQL · Forcing a broadcast
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 typeCan broadcast
innerEither side
left outer, left semi, left antiThe right side only
right outerThe left side only
full outerNeither, 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

Statistics can be wrong, especially after filters and joins. A 2 GB table forced into a broadcast must be collected to the driver and copied to every executor: driver out-of-memory or very slow jobs.

Too big for the hard limit

A broadcast table is limited to 8 GB; beyond that Spark fails the query.

Timeouts

Building a broadcast that takes longer than spark.sql.broadcastTimeout (300 seconds) fails the query, typically when the small side is itself an expensive query.

Memory on every executor

Each executor holds a copy. With 200 executors, a 500 MB table occupies 100 GB of cluster memory.
Tip: raising the threshold to 100 to 200 MB is a common, reasonable tuning step on clusters with plenty of memory. Raising it to gigabytes usually is not.

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

Spark ignores the hint and falls back to another strategy.

Trusting size estimates after complex transformations

Check actual sizes in the Spark UI.

Setting the threshold to several GB

Driver and executor memory pressure, and timeouts.

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 questions

1. 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.

Solve: Pick the Join Strategy →

Go deeper

Primary sources: Join strategy hints · Performance tuning