All concepts

Joins at Scale

If one side fits in memory, ship it to every machine and the join is free — otherwise both sides get shuffled, sorted, and merged.

Distributed Processing · Advanced · ~6 min

In plain English

Either you photocopy the small list and hand a copy to everyone, or both sides queue up in the same order and walk past each other.

Why it's worth your time

Turning a sort-merge join into a broadcast join is usually the single biggest win available in a Spark pipeline.

If you remember three things

  • Broadcast when one side fits in executor memory — no shuffle at all
  • Sort-merge shuffles both sides; correct at any size, expensive
  • Skewed or null keys serialise the whole stage behind one task

Overview

A distributed join has to get matching keys onto the same machine, and there are only two ways to do it. Broadcast: if one side is small, send a full copy to every executor and each task joins its own partition locally — no shuffle at all. Sort-merge: if both sides are large, hash-partition both by the join key, sort each partition, and merge — correct at any size, and expensive, because it shuffles both datasets. Every join tuning story is about turning the second into the first, or about making the second's partitions even.

In an interview

Broadcast hash join copies the small side to every executor and joins locally with no shuffle — the fastest option, limited by the small side fitting in executor memory. Sort-merge join shuffles and sorts both sides by the key, then merges; it works at any size and costs a full shuffle of both. Shuffle hash join sits between. The optimiser picks using size statistics, which are only as good as your table stats.

Production defaults

Statistics
keep them fresh — most 'wrong join strategy' tickets end at a missing ANALYZE
Order of work
filter and project both sides before the join, never after
Skew
adaptive skew join first; hand-salt only if the plan shows it wasn't enough

What breaks

  • Broadcast caused an OOM — The 'small' side wasn't. Check its actual size and raise or remove the hint.
  • Output row count exploded — A many-to-many join. Assert the expected count — silent row multiplication surfaces in a metric weeks later.

Watch it explained

What is Broadcast Join in spark? | Spark Optimization | IN 3 MINUTES | Definition | Applications — Quick Tech Bits, 4:34

Related