This should be a broadcast join, so the first thing is to check whether Spark actually did one. A 10 MB table joined with a 2 TB table taking hours almost always means Spark chose a sort-merge join and shuffled the 2 TB side.
Step 1: read the plan
joined.explain()
Look for BroadcastHashJoin. If you see SortMergeJoin with an Exchange on both sides, the broadcast did not happen. In the UI, the SQL tab shows the same plan with real sizes.
Step 2: why it did not broadcast
- The threshold was turned off.
spark.sql.autoBroadcastJoinThreshold = -1disables automatic broadcasts. Someone may have set it globally to avoid an old problem. - Spark's size estimate is wrong. The table is 10 MB on disk as compressed Parquet but Spark thinks it is large, because statistics are missing, or the lookup is the result of a complex subquery with no estimate. In that case, Spark does not know that it is small.
- The small side is only small after filtering, and Spark could not see that at planning time. AQE may fix this at runtime, if it is on.
Step 3: force it
from pyspark.sql.functions import broadcast fact.join(broadcast(lookup), "product_id")
Or the SQL hint /*+ BROADCAST(lookup) */. A hint overrides the estimate, so only use it when you know the side is small.
Limits that surprise people
- The broadcast side has to fit in driver memory (it is collected there first) and in every executor. "10 MB on disk" can be much larger in memory after decompression, but 10 MB is nowhere near a problem.
- For outer joins, only one side can be broadcast. In a left outer join, you can broadcast the right side but not the left, because the left side's unmatched rows need to be kept. If your small table is on the preserved side, broadcast cannot help, and Spark falls back to another strategy.
- Broadcast has a timeout (
spark.sql.broadcastTimeout, 300 seconds by default).
After the fix
Compare the plan and the run time. Also check for skew or a Python UDF in the same job, because a slow join is not always the only slow step.