Menu

PySpark interview question · Question 2 of 2

When would you use a broadcast join in Spark?

  • Medium
  • conceptual / optimization
  • ~8 min
  • High relevance
  • 2 min read
  • Updated Oct 2026

Short answer

Use a broadcast join when one side of the join is small enough to fit in memory on every executor. Spark copies the small table to each executor, so the large table is joined where it already sits and never has to be shuffled. It turns a costly sort-merge join into a much cheaper hash join, but broadcasting a table that is not truly small can exhaust driver or executor memory.

On this page
  1. Detailed explanation
  2. Example
  3. When it becomes dangerous
  4. Common mistakes

Detailed explanation

In a sort-merge join, Spark shuffles both tables so rows with the same key land on the same executor, then sorts and merges them. Shuffling a large table is expensive.

In a broadcast hash join, Spark sends a full copy of the small table to every executor and builds an in-memory hash table from it. Each partition of the large table is joined locally. The large side is not shuffled at all.

Example

from pyspark.sql import functions as F

orders = spark.read.parquet("/data/orders")       # large
countries = spark.read.parquet("/data/countries") # a few hundred rows

joined = orders.join(F.broadcast(countries), "country_code")
joined.explain()   # look for BroadcastHashJoin

Spark also broadcasts automatically when it estimates a side is below spark.sql.autoBroadcastJoinThreshold (10 MB by default), and Adaptive Query Execution can switch to a broadcast join at run time once it sees actual sizes.

When it becomes dangerous

  • The “small” table is not small. It must be collected and held in memory on the driver and every executor, which can cause out-of-memory errors.
  • Statistics are wrong. A filtered table may look small to the optimizer but be large in practice, or the reverse.
  • Outer-join direction. For a left outer join, Spark can only broadcast the right side; the side whose unmatched rows must be preserved cannot be the broadcast side.

Common mistakes

  1. Forcing a broadcast hint on a table that grows over time.
  2. Raising the threshold globally to “fix” one slow job.
  3. Not checking the physical plan to confirm which join actually ran.

By Data Career Hub Editorial · Last reviewed Oct 2026 · Defaults checked on Spark 4.2

Progress is saved in this browser only. No account needed.

Search
Filter by type