Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. 2 TB table joined to a 10 MB table is slow

PySpark · Joins & Data Layout

2 TB table joined to a 10 MB table is slow

Mediumpyspark-65
scenariobroadcast-joinexplainjoin-optimization

Question

A join between a 2 TB fact and a 10 MB lookup table takes hours. What do you check first?

Solution

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

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext