← All writing
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 autoBroadcastJoinThreshold is 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?

This is the kind of problem I get hired to fix. If it sounds like your pipeline, let's talk.

Message me on LinkedIn