spark/core
intermediate
Connecting…

Join strategies in Spark

Sort-merge shuffles both sides. Broadcast copies the small side. The wrong choice shuffles terabytes.

Lesson 14 of 29 · Spark path

Explain it at my level

  1. Scan
  2. Strategy
  3. Join
Watch the canvas:fact tabledimensionnetwork moveLive simulation
ordersfact · 8B rows2 TBregionsdimension · 12 rows20 MBSortMergeJoinboth sides must shuffle2TB20RExecutor 1idleRExecutor 2idleRExecutor 3idleCorrect — but 2 TB crossed the network for a 20 MB table4 GB × 3 executors = 12 GB of copies. Broadcast just became an OOM.
The cheapest join is the one where the big table never moves.