PySpark interview questionsQuestion 2 of 2
PySpark interview question · Question 2 of 2
When would you use a broadcast join in Spark?
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.
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
- Forcing a broadcast hint on a table that grows over time.
- Raising the threshold globally to “fix” one slow job.
- Not checking the physical plan to confirm which join actually ran.
Progress is saved in this browser only. No account needed.