Lesson 17 / 25

Join Strategies and Broadcast Joins

Pick broadcast hash join for small dimensions and understand the sort-merge default.

Avoid shuffling the big side

Spark chooses a physical join strategy. Broadcast hash join copies a small table (under spark.sql.autoBroadcastJoinThreshold, default 10 MB, or hinted with F.broadcast) to every executor so the large table is never shuffled: fastest, but the small side must fit in memory on the driver and every executor. Sort-merge join is the default for two large tables: both sides are shuffled and sorted by the join key. Shuffle hash join builds hash tables per partition. To speed up joins: filter and select before joining, broadcast small dimensions, join on integer keys where possible, handle skew, and avoid joining on expressions that block optimisation. With AQE, Spark can convert to a broadcast join at run time when a side turns out small.

Hinting a broadcast (illustrative)

Both forms ask Spark to broadcast dim_customers. The plan from the previous section showed the resulting BroadcastHashJoin.

from pyspark.sql import functions as F

result = fact_orders.join(F.broadcast(dim_customers), "customer_id", "left")

# SQL hint form
spark.sql("""
  SELECT /*+ BROADCAST(c) */ o.*, c.tier
  FROM fact_orders o LEFT JOIN dim_customers c ON o.customer_id = c.customer_id
""")

Do not broadcast something big

Broadcasting a table that does not fit in executor memory causes out-of-memory failures. Check its size, and prefer the default join when unsure.

Quick check: What is the main benefit of a broadcast join?

  • It works for any table size
  • It sorts both tables
  • The large table does not need to be shuffled
  • It deletes duplicates
Answer

The large table does not need to be shuffled — Shipping the small table everywhere replaces an expensive shuffle with a cheap local lookup.