Skip to content
Good engineers know 8 min read · spark.sql.join.preferSortMergeJoin 1 practice problem ↓

Join strategies

Broadcast hash, shuffle hash, sort-merge and nested loop joins: how Spark picks one, and how to steer it.

TL;DR The join type (inner, left, right, full, semi, anti, cross) decides which rows come back. The join strategy decides how Spark computes them: broadcast hash, shuffle hash, sort-merge, broadcast nested loop or Cartesian product. Large equi-joins default to sort-merge; joins without an equality condition can only use the two nested-loop strategies, which is why a BETWEEN join can run for hours.
Read first: Narrow vs wide transformations (skip if you know it)

Join types vs join strategies

Join typeReturns
innerOnly rows with a match on both sides
left / right outerEvery row of one side, with nulls where the other side has no match
full outerEvery row of both sides, matched where possible
left_semiLeft rows that have at least one match; left columns only, never duplicated
left_antiLeft rows that have no match
crossEvery combination of left and right rows

The type is part of what your query means. The strategy is an execution choice Spark makes during physical planning: the same inner join can run as any of the five strategies depending on table sizes, the join condition, configs and hints. The PySpark join and left_semi / left_anti lessons cover the types; this lesson covers the strategies.

The five strategies

Strategy (plan node)How it worksNeedsMain risk
Broadcastbroadcast join: A join where the small table is copied to every machine, so the big table never has to be shuffled. Learn more → hash join
BroadcastHashJoin
Copy the small side to every executorexecutor: A worker process on a cluster machine that runs tasks and holds data in memory. Learn more →, build a hash table from it, stream the large side through it. No shuffleshuffle: Moving rows between machines so that all rows with the same key end up together. Needed by joins, groupBy and sorting, and usually the most expensive step of a job. Learn more →.Equi-join; one side small; not a full outer joinDriverdriver: The process that runs your program, plans the work and sends tasks to executors. Learn more → and executor memory when the "small" side is not small
Shuffle hash join
ShuffledHashJoin
Shuffle both sides by key. In each partitionpartition: A chunk of a DataFrame's rows. Spark processes each partition as one task, so partitions decide how much work runs in parallel. Learn more →, build a hash table from the smaller side and probe it with the other. No sort.Equi-join; the build side of every partition fits in memoryOut-of-memory on one large or skewedskew: When one key has far more rows than the others, so one task does most of the work while the rest wait. Learn more → partition
Sort-merge join
SortMergeJoin
Shuffle both sides by key, sort each partition, then walk the two sorted streams together.Equi-join; sortable key typesTwo shuffles and two sorts; skewed keys
Broadcast nested loop join
BroadcastNestedLoopJoin
Broadcast one side and compare every row of the other side with every broadcast row.Any condition, or none; any join typerows × rows comparisons; slow when both sides are large
Cartesian product
CartesianProduct
Pair every partition of one side with every partition of the other, and test the condition on every row pair.Inner or cross joinWork and output grow as rows × rows

An equi-join has at least one equality between a column of each side in its condition (o.customer_id = c.customer_id, possibly ANDed with other conditions). Only equi-joins can use the hash and sort strategies, because those need a key to partition and match on. A condition with no such equality, like e.ts BETWEEN s.start AND s.end, a.x < b.y, or two equalities joined by OR, is a non-equi join and can only run as a nested loop.

How Spark picks one

For an equi-join, Spark walks this list and takes the first strategy that applies:

  1. 1
    Hints. If you hinted a strategy and the join supports it, Spark uses it. When both sides carry different hints, the priority is BROADCAST, then MERGE, then SHUFFLE_HASH, then SHUFFLE_REPLICATE_NL.
  2. 2
    Broadcast hash join if the join type allows broadcasting a side and that side's estimated size is below spark.sql.autoBroadcastJoinThreshold (10 MB). If both sides qualify, the smaller is broadcast.
  3. 3
    Shuffle hash join if spark.sql.join.preferSortMergeJoin is false (it is true by default), one side is small enough to hash per partition (below the broadcast threshold × the number of shuffle partitions), and it is at least three times smaller than the other side.
  4. 4
    Sort-merge join if the join keys are sortable. Most joins of two large tables end here.
  5. 5
    Cartesian product if it is an inner or cross join.
  6. 6
    Broadcast nested loop join as the last resort. It handles every case, even if that means broadcasting a large side.

For a non-equi join the list is shorter: a broadcast hint or a side below the threshold gives a broadcast nested loop join; otherwise an inner or cross join becomes a Cartesian product; otherwise Spark falls back to a broadcast nested loop join anyway.

Note: with AQEAQE: Adaptive Query Execution: Spark re-plans a running query using the real data sizes it has measured. Learn more → on (the default), the choice is revisited after each shuffle using measured sizes. A planned sort-merge join can become a broadcast hash join when one side turns out small, or a shuffle hash join when every partition is below spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold (0 by default, so off).

Why sort-merge is the default

A shuffle hash join skips the sort, so on paper it is faster. But its hash table for a partition must fit in memory, and one oversized or skewed partition fails the tasktask: The work for one partition in one stage, run on one CPU core. Learn more →. A sort-merge join sorts with an external sorter that spillsspill: Writing data to local disk because it does not fit in memory. Slower, but the job keeps running. Learn more → to disk, and the merge holds only the rows of the current key, so it slows down under pressure instead of failing. That robustness is why Spark prefers it. Shuffle hash join is worth trying when one side is too big to broadcast but small per partition, and keys are not skewed.

Which strategies support which join types

Join typeBroadcast hashShuffle hashSort-mergeBroadcast nested loopCartesian
inner, crossYes, either sideYesYesYes, either sideYes
left outer, left semi, left antiRight side broadcastYesYesRight side broadcastNo
right outerLeft side broadcastYesYesLeft side broadcastNo
full outerNoYes (Spark 3.1+)YesYes, but scans data several timesNo

The pattern: the side that must keep unmatched rows is streamed, so only the other side can be the broadcast or build side.

Hints

PySparkSpark SQL · Asking for a strategy
orders.join(customers.hint("broadcast"), "customer_id")      # BroadcastHashJoin
orders.join(customers.hint("merge"), "customer_id")          # SortMergeJoin
orders.join(customers.hint("shuffle_hash"), "customer_id")   # ShuffledHashJoin
events.join(sessions.hint("shuffle_replicate_nl"),
            events.ts.between(sessions.start, sessions.end))  # CartesianProduct
SELECT /*+ BROADCAST(c) */ * FROM orders o JOIN customers c ON o.customer_id = c.customer_id;
SELECT /*+ MERGE(c) */ * FROM orders o JOIN customers c ON o.customer_id = c.customer_id;
SELECT /*+ SHUFFLE_HASH(c) */ * FROM orders o JOIN customers c ON o.customer_id = c.customer_id;
SELECT /*+ SHUFFLE_REPLICATE_NL(s) */ * FROM events e JOIN sessions s ON e.ts BETWEEN s.start AND s.end;

A hint is a request, not an order. Spark ignores a hint it cannot honour, such as broadcasting the preserved side of an outer join or a shuffle hash join on a non-equi condition, and logs a warning. MERGEJOIN and SHUFFLE_MERGE are aliases of MERGE.

Making a slow non-equi join fast

  • OR across two keys — a.id = b.id OR a.email = b.email has no single equality, so it becomes a nested loop. Rewrite it as two equi-joins and union the results (deduplicating if a row can match both ways).
  • Range joins — Matching an IP to an IP range, or an event to a time window, has no equality. Add one the condition already implies: a coarse bucket (the /16 prefix of the IP, the hour of the timestamp). Explode each range into the buckets it covers, join on bucket equality, and keep the range test as an extra condition. The join becomes a sort-merge join on the bucket.
  • Functions on the key — A join on lower(a.code) = b.code is still an equi-join, but nothing about it can use bucketing or existing partitioning. Normalise the key once, before the join.

Check what Spark chose

Search the physical plan for the join node. SortMergeJoin with an Exchange and Sort under both sides is the expensive default; BroadcastHashJoin with a BroadcastExchange under one side avoids the big shuffle; BroadcastNestedLoopJoin or CartesianProduct on two large inputs is the red flag. With AQE, check the final plan in the SQL tab of the Spark UI rather than the plan before execution.

Common mistakes

  • Confusing join type with join strategy — A left join is not a strategy. It can run as a broadcast hash, shuffle hash, sort-merge or nested loop join.
  • Forcing shuffle hash joins on skewed data — One hot key builds a hash table that does not fit, and the task fails instead of spilling.
  • Not noticing a nested loop — A non-equi condition on two large tables quietly turns into rows × rows comparisons.
  • Expecting a hint to always win — Spark drops hints that the join type or condition cannot support.

What you learned

  • The difference between a join type and a join strategy
  • The five physical join strategies and what each needs
  • The order in which Spark picks a strategy
  • How hints steer the choice, and how to fix a slow non-equi join

Key takeaways

  • Join type is what you get; join strategy is how Spark computes it.
  • Equi-joins can use broadcast hash, shuffle hash or sort-merge; non-equi joins only nested loops.
  • Order: hints, broadcast hash, shuffle hash (if preferred), sort-merge, Cartesian, broadcast nested loop.
  • Sort-merge is the default for large joins because it spills rather than failing.
  • Rewrite OR and range conditions into equi-joins to escape nested loops.

Check yourself

4 questions

Which strategy does a join on e.ts BETWEEN s.start AND s.end use when neither side is small?

Show the answer

CartesianProduct, for an inner join. There is no equality, so only nested-loop strategies apply; an inner join of two large sides becomes a Cartesian product.

Why does Spark prefer sort-merge over shuffle hash join by default?

Show the answer

It spills to disk under memory pressure, while a per-partition hash table must fit in memory. spark.sql.join.preferSortMergeJoin is true because sort-merge degrades gracefully.

Which join type can never use a broadcast hash join?

Show the answer

full outer. Both sides must keep unmatched rows, so neither can be the broadcast side.

Both sides of a join are hinted, one MERGE and one SHUFFLE_HASH. What wins?

Show the answer

MERGE. The priority is BROADCAST, MERGE, SHUFFLE_HASH, SHUFFLE_REPLICATE_NL.

Practice it

Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.

Solve: Pick the Join Strategy →

Keep going

Up next · lesson 8 of 30 · 3 min read
Broadcast joins
Ship a small table to every executor and skip the shuffle, and when broadcasting backfires.

Related lessons

Previous: Narrow vs wide transformations

Primary sources: Join strategy hints · Performance tuning: join strategy hints