A broadcast (map-side) join copies the small side to every executor, then joins locally. It avoids shuffling the large fact table.
Diagram
Small dim ---broadcast---> Executor1 \
Executor2 +-- local join with local fact partitions
Large fact --no shuffle--> Executor3 /Code
from pyspark.sql.functions import broadcast
fact = spark.read.parquet("/lake/fact_orders") # huge
dim = spark.read.parquet("/lake/dim_users") # small
out = fact.join(broadcast(dim), "user_id", "inner")When it helps
- Dimension tables, lookup codes, small filter lists
- Below broadcast threshold (default often 10MB, configurable via
spark.sql.autoBroadcastJoinThreshold)
When it hurts / fails
- "Small" table is actually large -> driver/executor OOM
- Streaming unstable state sizes
- Threshold too aggressive for skewed clusters
Interview tip
Mention AQE can convert sort-merge joins to broadcast at runtime when size estimates warrant it, and that you still hint when you know the dimension is tiny.