Skip to content
Reliable Data Engineering
Lesson
Open in the interactive app

Spark Join Strategies and Broadcast Joins

Joins are where most Spark time and most interview questions go. This module explains how a logical join becomes a physical strategy, how to read and influence that choice, how broadcast joins really work (and fail), and what AQE changes at runtime.


1. The five physical strategies

StrategyHow it worksNeedsCost profileSupports
Broadcast hash join (BHJ)Collect the small side to the driver, ship it to every executor, build a hash table; stream the big side through itEqui-join; small side fits in memory (≤ 8 GB hard limit, practically ≤ a few hundred MB)No shuffle of the big side: fastest for star joinsInner, left/right outer (broadcast the non-preserved side), semi, anti; not full outer
Shuffle hash join (SHJ)Shuffle both sides by key; in each partition build a hash table from the smaller sideEqui-join; build side per partition fits in memoryOne shuffle, no sort; risky if a build partition is bigMost join types (full outer since 3.1)
Sort-merge join (SMJ)Shuffle both sides by key, sort each partition, mergeEqui-join with sortable keysShuffle + sort of both sides; robust because it spillsAll equi-join types; the default for large-large
Broadcast nested loop join (BNLJ)Broadcast one side; compare every pair with the conditionAny condition (non-equi)O(n·m) comparisons; OK only if one side is tinyAll types
Cartesian productEvery row × every rowExplicit cross join / no conditionO(n·m) outputInner only
flowchart TD
    J[Logical join] --> EQ{Has equality keys?}
    EQ -->|yes| H{Hint present?}
    H -->|"BROADCAST / MERGE / SHUFFLE_HASH / SHUFFLE_REPLICATE_NL"| HINT[Use the hinted strategy if valid for the join type]
    H -->|no| B{One side ≤ autoBroadcastJoinThreshold<br/>and can be broadcast for this join type?}
    B -->|yes| BHJ[Broadcast hash join]
    B -->|no| S{"preferSortMergeJoin = false<br/>and one side much smaller and<br/>fits per-partition hash map?"}
    S -->|yes| SHJ[Shuffle hash join]
    S -->|no| SMJ[Sort-merge join]
    EQ -->|no| NB{One side broadcastable?}
    NB -->|yes| BNLJ[Broadcast nested loop join]
    NB -->|no, inner| CP[Cartesian product]
    NB -->|no, outer| BNLJ2[BNLJ anyway: risky]

Notes on the selection logic:


2. Join type and strategy compatibility

Join typeBHJSHJSMJBNLJ
Innereither sideeither side✓✓
Left outer / left semi / left antibroadcast rightbuild right✓broadcast right
Right outerbroadcast leftbuild left✓broadcast left
Full outer✗✓ (3.1+)✓✓ (slow)
Cross✗ (equi only)✗✗✓ / Cartesian

Interview trap: “Why didn’t my left join broadcast the small left table?” Because in a left outer join the left side is preserved and can’t be the broadcast side. Broadcast the right side, or rewrite.


3. Statistics, size estimation and the cost-based optimizer


4. Join hints

HintEffect
BROADCAST / BROADCASTJOIN / MAPJOINBroadcast this side (if allowed for the join type), regardless of the threshold
MERGE / SHUFFLE_MERGE / MERGEJOINUse sort-merge join
SHUFFLE_HASHUse shuffle hash join with this side as the build side
SHUFFLE_REPLICATE_NLUse a Cartesian/nested-loop strategy
SELECT /*+ BROADCAST(d) */ f.*, d.region FROM sales f JOIN stores d ON f.store_id = d.store_id;
SELECT /*+ SHUFFLE_HASH(o) */ * FROM orders o JOIN payments p USING (order_id);
from pyspark.sql.functions import broadcast
fact.join(broadcast(dim), "store_id")
orders.hint("shuffle_hash").join(payments, "order_id")

5. Broadcast joins in depth

How a broadcast join runs

sequenceDiagram
    participant D as Driver
    participant E1 as Executor 1
    participant E2 as Executor 2
    D->>E1: run job to compute the small side
    D->>E2: (tasks in parallel)
    E1-->>D: collect partitions of small side
    E2-->>D: collect partitions of small side
    Note over D: build HashedRelation<br/>(serialized, compressed)
    D->>E1: torrent-style broadcast (chunks)
    D->>E2: executors also share chunks peer-to-peer
    Note over E1,E2: each executor deserializes ONE copy<br/>shared by all its tasks
    E1->>E1: stream big-side partitions through the hash table
    E2->>E2: no shuffle of the big side
  1. The small side is computed and collected to the driver (so driver memory must hold it).
  2. The driver builds a HashedRelation (a hash map keyed by the join key), serializes and compresses it, and splits it into chunks.
  3. TorrentBroadcast distributes the chunks; executors fetch from the driver and from each other.
  4. Each executor deserializes one copy into memory (storage memory), shared by all tasks on it.
  5. Big-side tasks probe the hash table locally: no shuffle, no sort of the big side.

Sizing and limits

Failure modes

SymptomCauseFix
Driver OOM during a joinA large table broadcast (wrong stats, aggressive hint/threshold)Fix stats, remove the hint, lower the threshold; let AQE decide
Could not execute broadcast in 300 secsSmall side is expensive to compute (big upstream job) or too largeMaterialise or cache the small side first, filter earlier, raise the timeout only if the size is fine
Executor OOM / heavy GC in BHJ stagesThe hash relation is big in deserialized form × the executor’s memoryLower the threshold or use SHJ/SMJ; more memory per executor
Missed broadcast (SMJ on a small table)Estimate is too big (no stats, filter selectivity unknown)ANALYZE statistics, explicit broadcast() hint, or rely on AQE runtime conversion
Broadcast reused wrongly after cache()A cached DataFrame’s size estimate replaces the source estimateCheck explain() after caching (see the caching module)

Broadcast join vs broadcast variable

A broadcast join is planned by Spark SQL for DataFrames. A broadcast variable (sc.broadcast(obj)) ships any read-only object (a lookup dict, a model) to executors once instead of capturing it in every task closure. Use joins for tabular lookups and variables for small non-tabular objects in UDFs and RDD code.


6. AQE: runtime join changes

With adaptive execution, Spark re-plans after each shuffle map stage using actual sizes:

AQE ruleWhat it does
SMJ → BHJIf a side’s actual shuffle output is ≤ spark.sql.adaptive.autoBroadcastJoinThreshold (defaults to the static threshold), convert to broadcast; the local shuffle reader then reads already-written shuffle files locally instead of fetching over the network
SMJ → SHJIf all partitions of a side fit under spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold, use shuffle hash join (no sort)
Demote broadcastDon’t broadcast a side whose shuffle output is mostly empty partitions (non-empty ratio below spark.sql.adaptive.nonEmptyPartitionRatioForBroadcastJoin), where a broadcast would be inefficient (DemoteBroadcastHashJoin)
Skew joinSplit skewed partitions of an SMJ/SHJ and replicate the matching side (see skew)
Coalesce partitionsMerge small post-shuffle partitions before the join’s reduce tasks

Read the final plan after the query runs (AdaptiveSparkPlan isFinalPlan=true in the SQL tab) to see what actually happened. The initial plan can say SMJ while the executed plan used BHJ.


7. Non-equi joins, range joins and tricky conditions

-- events joined to sessions by time range: bucket by hour (a session can span hours, so explode its hours)
WITH s AS (SELECT *, explode(sequence(date_trunc('hour', start_ts), date_trunc('hour', end_ts), INTERVAL 1 HOUR)) AS hr FROM sessions)
SELECT e.*, s.session_id
FROM events e JOIN s
  ON e.user_id = s.user_id AND date_trunc('hour', e.ts) = s.hr      -- equi part → hash/sort strategies
 AND e.ts BETWEEN s.start_ts AND s.end_ts                           -- precise part

Databricks also provides a range join optimization hint (/*+ RANGE_JOIN(s, 3600) */) that does this binning internally.


8. Runtime filters, DPP and storage-partitioned joins


9. Optimisation playbook and decision framework

  1. Shrink before joining: filter, project, and aggregate inputs; semi-join to reduce one side (WHERE EXISTS / left_semi).
  2. Small side ≤ a few hundred MB? Broadcast it (fix stats or hint); make sure the join type allows it.
  3. Both large? SMJ by default; consider SHJ if one side is moderately small per partition and sorting dominates.
  4. Repeated large-large joins on the same key? Co-partition: Iceberg SPJ or bucketed Hive tables, or persist both sides repartitioned by the key within one job.
  5. Skewed keys? AQE skew join → salting hot keys → handle NULL/default keys separately.
  6. Non-equi condition? Add an equi bucket key (binning), use range join hints, or restructure.
  7. Many joins? Fresh statistics + CBO join reordering, or reorder manually (selective joins first).
  8. Verify in the executed plan and the SQL tab: strategy per join, rows in/out, shuffle sizes, spill.

Override the optimizer when: you know a side is tiny but the stats say otherwise (broadcast()), you know a broadcast is dangerous as data grows (MERGE), or a known build side should be enforced (SHUFFLE_HASH). Document why next to the hint.


Interview questions

Walk me through how Spark chooses a join strategy for an equi-join.

Hints first (BROADCAST > MERGE > SHUFFLE_HASH > SHUFFLE_REPLICATE_NL, if valid for the join type). Otherwise, if one side’s estimated size is under autoBroadcastJoinThreshold and that side can be the build side for the join type, broadcast hash join. Otherwise, if sort-merge isn’t preferred and one side is small enough for per-partition hash maps, shuffle hash join; else sort-merge join. Without equality keys: broadcast nested loop join or a Cartesian product. AQE can then change the strategy at runtime using actual sizes.

Why might Spark broadcast a huge table, and how do you prevent it?

Because the decision uses size estimates: stale or missing statistics, compressed file sizes that understate memory size, or a hint/threshold set too high. Symptoms are driver OOM or a broadcast timeout. Fix the statistics (ANALYZE TABLE), remove or lower the hint/threshold, and rely on AQE’s runtime conversion, which uses actual shuffle sizes.

A LEFT JOIN with a small left table isn't broadcast. Why?

In a left outer join, the left side is preserved (all its rows must appear), so only the right side can be the broadcast/build side. If the right side is big, Spark falls back to sort-merge. Rewrite (e.g. swap to a right join or an inner join plus a union for unmatched rows), or accept SMJ.

How do you join events to sessions where event time falls between session start and end, at scale?

A pure range condition forces a nested loop. Add an equality on user_id plus a time bucket: explode each session into the hourly (or daily) buckets it spans, join on (user_id, bucket), then apply the exact BETWEEN filter. On Databricks the RANGE_JOIN hint does this binning automatically. Pick the bucket size near typical session length to balance replication and selectivity.

What does AQE change about joins at runtime?

After shuffle map stages finish, it uses real sizes to convert sort-merge to broadcast (then reads shuffle files locally), convert sort-merge to shuffle hash join when partitions are small, avoid inefficient broadcasts, split skewed partitions for skew joins, and coalesce small partitions. The initial plan in explain() can differ from the final executed plan.

What's a storage-partitioned join?

A shuffle-free join between V2 tables (e.g. Iceberg) whose storage partitioning on the join key is compatible (such as matching bucket transforms). Spark uses the reported partitioning to join corresponding partitions directly. It’s the modern table-format equivalent of bucketed joins in Hive tables, and is enabled via the v2 bucketing configs.