Spark
Spark is silently killing your join performance
Spark is silently killing your join performance, and the fix is one config change.
Most data engineers know about broadcast joins. Few actually have them working correctly in production.
Here’s the trap:
- Spark’s default
autoBroadcastJoinThresholdis 10MB - Your dimension tables are 50MB, 100MB, 200MB
- Spark looks at them, decides they are too big to broadcast, and falls back to a shuffle join
That means:
- Gigabytes of data moving across the network
- Executors waiting on each other
- Your 5-minute job running for 40 minutes
The fix:
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", 209715200)
One line. Done.
Know your dimension table sizes and set the threshold accordingly. Most engineers set up Spark with defaults and never touch them again. The defaults are safe. They’re not optimal.
What default config have you changed that made the biggest difference in your pipelines?