Spark
Mid-levelWhen does Spark's broadcast join actually help, and when does it backfire?
A broadcast join is a bet that one side of the join is small enough to copy everywhere cheaply. When the bet is right it's the fastest join Spark has. When it's wrong, it's the fastest way to OOM a driver.
spark.sql.autoBroadcastJoinThreshold (10MB by default) tells the planner: if a table's estimated size is under this, skip the shuffle and send a full copy to every executor instead.
What the threshold actually controls
Below the threshold, the planner replaces a shuffle join (both sides repartitioned and shuffled across the network) with a broadcast join (the small side collected to the driver, then pushed to every executor as a hash table). No shuffle for the large side means no shuffle-read stage, no shuffle-related skew, and usually a meaningfully faster join - which is exactly why the planner does this automatically whenever it estimates a table is small enough.
The failure mode: a broadcast that should have stayed a shuffle
The estimate the planner uses comes from table statistics, which go stale - a table that was 8MB when stats were last computed but is now 800MB will still get broadcast if nothing refreshed those stats. The result: every executor tries to hold an 800MB hash table in memory at once, and the driver (which collects the broadcast data first) is often the one that OOMs before any executor does.
How to tell which one you're looking at from the event log
Look for a stage with very few tasks (a broadcast collection stage runs on essentially one task) showing unusually high peak execution memory, immediately followed by a stage-level failure or a driver OOM in the application logs rather than a task-level one. A shuffle-join OOM, by contrast, shows up as many tasks failing across a normal-sized stage - a different signature entirely.
If you're not sure a broadcast is happening at all: check the physical plan for
BroadcastHashJoin vs SortMergeJoin - it's the one place Spark tells you
outright which strategy it picked, rather than leaving you to infer it from timing.
See what this looks like on your own pipeline.
Upload a real Spark event log, dbt run_results.json, or Flink metrics export and get
concrete, approve-before-apply recommendations back - not another rule of thumb.