All posts

Apache Spark joins in depth: the five strategies and how one gets picked

Spark has five join strategies and one decision procedure that chooses between them. Here is how each works, the exact order Spark tries them in, which join types quietly forbid which strategy, and what to do per scenario.

33 min read Spark

TL;DR

  • A join type is semantics (inner, outer, semi, anti). A join strategy is the algorithm (broadcast hash, shuffle hash, sort-merge, broadcast nested loop, cartesian). Mixing the two words up is why join tuning feels arbitrary.
  • Spark’s choice is a documented, ordered procedure in JoinSelection, not a cost model. Hints first, then broadcast, then shuffle hash, then sort-merge, then cartesian, then nested loop “as the final solution”.
  • The join type silently restricts the strategy. A full outer join can never be broadcast, and for a right outer join only the left side may be. Both are confirmed below by running them.
  • Shuffle hash join needs three conditions at once, one of which is off by default, which is why you almost never see it without a hint.
  • The 10 MB threshold is not the only size rule. Shuffle hash join has its own ceiling of autoBroadcastJoinThreshold multiplied by spark.sql.shuffle.partitions, which is 2 GB on the defaults, and a side effect: setting that threshold to -1 disables shuffle hash join along with broadcasting.

Join tuning has a reputation for being folklore. Broadcast the small side, raise the threshold, add a hint, and if none of that works, salt the key. It works often enough to feel like knowledge and fails often enough to feel like luck.

The reason is that most of what circulates describes the five strategies without describing the procedure that picks between them. That procedure is not mysterious: it is written down, in order, in a comment above JoinSelection in Spark’s own source, and once you have read it the behaviour stops being surprising. A full outer join that refuses to broadcast is not a bug, and no amount of raising the threshold will change it.

This post has three parts. What each strategy actually does, the exact order Spark tries them in, and which scenario should get which.

What is the difference between a join type and a join strategy?

This is the distinction to get right first, because the two are chosen independently and the first one constrains the second.

  Join type Join strategy
Decides Which rows appear in the result How the rows are brought together
You write it Yes, it is your query’s meaning No, Catalyst picks it
Examples inner, left outer, full outer, left semi, left anti, cross broadcast hash, shuffle hash, sort-merge, broadcast nested loop, cartesian
Changes results Yes Never

A strategy change can make a query ten times faster or fail outright, and it cannot change the answer. A type change always changes the answer. So when someone says “we changed the join and it got faster”, the useful question is which of the two they changed.

The types, briefly, since the rest of the post assumes them:

trips.join(cities, "city_id")                      # inner, the default
trips.join(cities, "city_id", "left")              # left outer
trips.join(cities, "city_id", "full")              # full outer
trips.join(cities, "city_id", "left_semi")         # left rows that match
trips.join(cities, "city_id", "left_anti")         # left rows that do not
trips.crossJoin(cities)                            # every pair

Semi and anti deserve a word because they are under-used. Both filter the left side by whether a match exists and return no columns from the right, so a left row matching a hundred right rows is emitted once. An inner join in the same position emits it a hundred times and needs a distinct afterwards, which is a second shuffle.

Architecture: what every join is made of

All five strategies are assembled from the same three parts, and naming them makes the differences between the strategies easy to hold.

flowchart TB
  subgraph P["every join has these"]
    A["<b>build side</b><br/>materialised into a lookup<br/>structure, or sorted"]
    B["<b>stream side</b><br/>read once and probed<br/>against the build side"]
    C["<b>co-location step</b><br/>how matching keys are made<br/>to meet on one executor"]
  end
  A --> J["the join operator"]
  B --> J
  C --> J

The build side is the one Spark materialises: a hash table for the hash joins, a sorted run for sort-merge. It is the side that costs memory, which is why “which side is the build side” is the question behind most join failures.

The stream side is read once and probed. It costs nothing to hold, which is why you want the large table here.

The co-location step is where the strategies genuinely differ, and it is the expensive part:

Strategy How keys are co-located Network cost
Broadcast hash Copy the whole build side everywhere One copy per executor
Shuffle hash Shuffle both sides by key Full shuffle
Sort-merge Shuffle both sides by key, then sort Full shuffle plus sort
Broadcast nested loop Copy one side everywhere, compare all pairs One copy per executor
Cartesian Replicate partitions across each other Every partition pair

Read down that column and the whole subject compresses to one trade: either you copy the small side to where the data already is, or you move both sides to a common place. Broadcasting does the first, the shuffle-based strategies do the second, and bucketing is how you pay for the second once at write time instead of on every read.

Why are there five strategies, and not six?

The five strategies are not five unrelated algorithms. They are combinations of two independent choices, and separating them is the clearest way to hold the whole subject.

Choice one, the data exchange: how matching keys are made to meet on one executor.

  • Broadcast. Copy one side to every executor. The other side never moves.
  • Shuffle. Repartition both sides by the join key so matching keys land in the same partition.

Choice two, the join algorithm: what happens once the rows are together.

  • Hash. Build a hash table from one side, probe it with the other.
  • Sort-merge. Sort both sides, walk them with a cursor each.
  • Nested loop. Compare every pair, evaluating an arbitrary condition.

Two exchanges times three algorithms is six combinations. Spark ships five physical operators, and the grid shows exactly which cell is missing:

  Hash Sort-merge Nested loop
Broadcast BroadcastHashJoinExec does not exist BroadcastNestedLoopJoinExec
Shuffle ShuffledHashJoinExec SortMergeJoinExec CartesianProductExec

Why there is no broadcast sort-merge join is worth a moment, because it explains what sorting is actually for. Sorting buys bounded memory: you never hold a whole side, and the sort can spill. But if a side has already been copied to every executor, it is in memory by definition, so a hash table over it costs nothing extra and gives constant-time probes. Sorting both sides to get the same answer more slowly has no upside. Sorting only earns its place when you had to shuffle anyway and cannot afford to hold a side in memory.

So the decision Spark makes is really two decisions:

flowchart TB
  A["can one side be copied<br/>to every executor?"]
  A -->|"yes"| B{"equi-join?"}
  A -->|"no"| C{"equi-join?"}
  B -->|"yes"| D["BroadcastHashJoin"]
  B -->|"no"| E["BroadcastNestedLoopJoin"]
  C -->|"yes"| F{"hold a side in memory<br/>per partition?"}
  C -->|"no"| G["CartesianProduct<br/><i>inner-like only</i>"]
  F -->|"yes"| H["ShuffledHashJoin"]
  F -->|"no"| I["SortMergeJoin"]

Read it as: the exchange is decided by size, and the algorithm by the condition and by how much memory you can hold. An equi-join can hash or sort-merge; a non-equi join has only nested loop, which is why it collapses to two options. That single sentence predicts most of what the planner does.

What are the five join strategies?

Spark’s source documents all five with their limits. The table below is that comment, condensed:

Strategy Equi-join only Keys must be sortable Join types
Broadcast hash join Yes No All except full outer
Shuffle hash join Yes No All
Sort-merge join Yes Yes All
Broadcast nested loop No No All, but efficient only in some combinations
Cartesian product No No Inner-like only

An equi-join is one whose condition is equality on keys, a.id = b.id. Anything else, a range or an inequality, is a non-equi join, and only the bottom two rows can execute it at all. That single fact explains most surprises in the plan.

Broadcast hash join

Spark collects the small side to the driver, sends a copy to every executor, and each executor builds a hash table from it. Each partition of the large side then probes that table locally. The large side is never shuffled.

flowchart TB
  S["small side: cities"] --> D["driver collects it"]
  D --> BC[("broadcast:<br/>one copy per executor")]
  BC --> E1["executor 1<br/>hash table"]
  BC --> E2["executor 2<br/>hash table"]
  BC --> E3["executor 3<br/>hash table"]
  T1["trips partition 0<br/>does not move"] --> E1
  T2["trips partition 1<br/>does not move"] --> E2
  T3["trips partition 2<br/>does not move"] --> E3
  E1 --> R["output, no shuffle<br/>of the large side"]
  E2 --> R
  E3 --> R

Cost: one copy of the small side per executor, held in memory for the duration, plus the driver having to assemble it first.

That driver step is the part people forget. A broadcast that is too large fails on the driver, not the executors, and the error often arrives as a spark.sql.broadcastTimeout after the default 300 seconds rather than as an obvious out-of-memory.

Key points

  • Build side: the small one, if canBroadcastBySize passes and the join type allows that side.
  • Network: one copy of the small side per executor. The large side never moves.
  • Memory: the hash table sits on every executor, and on the driver first.
  • Wins when: one side is genuinely small after filters and projections are applied.
  • Fails when: the small side is not small, or the join type forbids that build side.
  • In the plan: BroadcastHashJoin with a BroadcastExchange child.

Shuffle hash join

Both sides are shuffled on the join key. Then, for each partition, Spark builds a hash table from the smaller side and probes it with the larger. No sorting.

flowchart TB
  A["riders<br/>the smaller side"] --> SH[("shuffle both sides<br/>by the join key")]
  B["trips<br/>the larger side"] --> SH
  SH --> P0["partition 0<br/>build hash from riders<br/>probe with trips"]
  SH --> P1["partition 1<br/>build hash from riders<br/>probe with trips"]
  SH --> P2["partition 2<br/>build hash from riders<br/>probe with trips"]
  P0 --> O["output, never sorted"]
  P1 --> O
  P2 --> O

It sits in a narrow band: one side too big to broadcast, but small enough that a per-partition hash table fits. It skips the sort that sort-merge pays, and in exchange it can run out of memory where sort-merge would spill.

Key points

  • Build side: the smaller side, per partition, not the whole table.
  • Network: a full shuffle of both sides, the same as sort-merge.
  • Memory: one hash table per partition, which is what can blow up.
  • Wins when: the build side is much smaller but still too big to broadcast, and the sort is the thing you want to skip.
  • Fails when: a partition’s build side does not fit, since a hash table cannot spill the way a sort can.
  • In the plan: ShuffledHashJoin with an Exchange hashpartitioning under each side and no Sort node.

Sort-merge join

Both sides are shuffled on the key, each partition is sorted, and the two sorted streams are walked together with a cursor on each.

flowchart TB
  A["trips"] --> S[("shuffle both sides<br/>by the join key")]
  B["cities"] --> S
  S --> S0["partition 0<br/>sort both sides"]
  S --> S1["partition 1<br/>sort both sides"]
  S0 --> M0["merge with one<br/>cursor per side"]
  S1 --> M1["merge with one<br/>cursor per side"]
  M0 --> O["output, sorted by key"]
  M1 --> O

The merge itself is the cheap part. Each side has a cursor, and whichever cursor points at the smaller key advances:

left  (sorted):  a a b c d ...
right (sorted):  a b b d e ...
                 ^ ^

Its virtue is bounded memory. Only the merge buffers are held, not a whole hash table, and the sort can spill, so a partition larger than memory still finishes. That is why it is the default for large-to-large joins, and why spark.sql.join.preferSortMergeJoin defaults to true.

The cost is the sort, and the requirement is that the keys are sortable. It is the only strategy with that constraint.

Key points

  • Build side: neither, in the hash-table sense. Both sides are sorted and streamed.
  • Network: a full shuffle of both sides, plus a sort per partition.
  • Memory: bounded, and the sort spills to disk rather than failing.
  • Wins when: both sides are large, which is why it is the default.
  • Fails when: rarely, and not on memory. Its failure mode is slowness from the shuffle and the sort, or one skewed partition, rather than an out-of-memory.
  • In the plan: SortMergeJoin with a Sort node above an Exchange on each side.

Broadcast nested loop join

One side is broadcast, and then every row of one side is compared against every row of the other, evaluating an arbitrary condition. It is the fallback that makes non-equi joins possible at all.

The source notes it is optimised for three specific shapes, and slow otherwise:

  1. Broadcasting the left side in a right outer join.
  2. Broadcasting the right side in a left outer, left semi, left anti or existence join.
  3. Broadcasting either side in an inner-like join.
flowchart TB
  S["one side, broadcast whole"] --> E1["executor 1"]
  S --> E2["executor 2"]
  P1["partition of the other side"] --> E1
  P2["partition of the other side"] --> E2
  E1 --> C1["compare every local row<br/>against every broadcast row"]
  E2 --> C2["compare every local row<br/>against every broadcast row"]
  C1 --> O["rows where the<br/>condition holds"]
  C2 --> O

Outside those, “we need to scan the data multiple times, which can be rather slow”. Seeing BroadcastNestedLoopJoin on two large inputs means the query is doing something close to a cross product.

Key points

  • Build side: the broadcast one, and only the shapes listed above are efficient.
  • Network: one copy of the broadcast side per executor.
  • Cost: the row comparisons, which grow as the product of the two row counts.
  • Wins when: the condition is not an equality and one side is genuinely tiny.
  • Fails when: both sides are large, where it is the “may OOM” last resort.
  • In the plan: BroadcastNestedLoopJoin with the condition printed inline.

Cartesian product

Also called shuffle-and-replicate nested loop. Every partition of one side is paired with every partition of the other, producing the full cross product. It supports inner-like joins only, and it is what you get from a deliberate crossJoin once broadcasting is off the table.

flowchart TB
  A0["A partition 0"] --> X["every A partition paired<br/>with every B partition"]
  A1["A partition 1"] --> X
  B0["B partition 0"] --> X
  B1["B partition 1"] --> X
  X --> P00["A0 x B0"]
  X --> P01["A0 x B1"]
  X --> P10["A1 x B0"]
  X --> P11["A1 x B1"]

With two partitions a side that is four tasks. With 200 a side it is 40,000, and that multiplication is the whole reason this strategy is a warning sign.

Key points

  • Build side: neither. Partitions are replicated against each other.
  • Network: every partition pair, so the shuffle grows as the product of the partition counts.
  • Cost: the output itself, which is the product of the two row counts.
  • Wins when: you actually want a cross product, on small inputs.
  • Fails when: you did not want one, which is the usual case.
  • In the plan: CartesianProduct, almost always worth investigating.

How does Spark actually choose?

Here is the part that turns folklore into a procedure. JoinSelection follows a fixed order, and the first applicable rule wins.

For an equi-join, hints are consulted first, in this order:

  1. BROADCAST: pick broadcast hash join if the join type supports it.
  2. MERGE: pick sort-merge join if the keys are sortable.
  3. SHUFFLE_HASH: pick shuffle hash join if the join type supports it.
  4. SHUFFLE_REPLICATE_NL: pick cartesian product if the join type is inner-like.

With no hints, for an equi-join:

  1. Broadcast hash join, if one side is small enough to broadcast and the join type allows that side to be the build side.
  2. Shuffle hash join, if one side can build a local hash map, is much smaller than the other, and spark.sql.join.preferSortMergeJoin is false.
  3. Sort-merge join, if the join keys are sortable.
  4. Cartesian product, if the join type is inner-like.
  5. Broadcast nested loop join as the last resort. The source is blunt about it: “It may OOM but we don’t have other choice.”

For a non-equi join the list is much shorter, which is why these are so often slow:

  1. Broadcast nested loop join, if one side is small enough to broadcast.
  2. Cartesian product, if the join type is inner-like.
  3. Broadcast nested loop join anyway.
flowchart TB
  H{"a join hint?"} -->|"yes"| HH["honour it if the type<br/>and keys allow"]
  H -->|"no"| EQ{"equi-join?"}
  EQ -->|"no"| NE["broadcast nested loop if a side is small,<br/>else cartesian if inner-like,<br/>else nested loop anyway"]
  EQ -->|"yes"| B{"a side small enough,<br/>and allowed as build side?"}
  B -->|"yes"| BHJ["BroadcastHashJoin"]
  B -->|"no"| SH{"local hash map fits,<br/>much smaller,<br/>preferSortMergeJoin false?"}
  SH -->|"yes"| SHJ["ShuffledHashJoin"]
  SH -->|"no"| SM{"keys sortable?"}
  SM -->|"yes"| SMJ["SortMergeJoin"]
  SM -->|"no"| C["cartesian if inner-like,<br/>else nested loop"]

The size rules, exactly

Three predicates decide “small enough”, and only the first is widely known.

Broadcast. canBroadcastBySize is sizeInBytes >= 0 && sizeInBytes <= autoBroadcastJoinThreshold, where the threshold is spark.sql.autoBroadcastJoinThreshold, default 10485760, which is 10 MB.

There is a detail worth having: when the statistics are runtime statistics, meaning adaptive execution has measured a completed stage, Spark uses spark.sql.adaptive.autoBroadcastJoinThreshold instead, falling back to the static one when that is unset. So plan-time and runtime broadcasting can be governed by two different numbers.

Shuffle hash, ceiling. canBuildLocalHashMapBySize is sizeInBytes < autoBroadcastJoinThreshold * numShufflePartitions. On the defaults that is 10 MB multiplied by 200:

10,485,760 x 200 = 2,097,152,000 bytes = ~2 GB

So the shuffle-hash ceiling is roughly 2 GB by default, two hundred times the broadcast threshold, and it moves when you change the shuffle partition count.

Shuffle hash, ratio. muchSmaller is a.size * spark.sql.shuffledHashJoinFactor <= b.size, and that factor defaults to 3. The source explains why: “The cost to build hash map is higher than sorting, we should only build hash map on a table that is much smaller than other one.”

So shuffle hash join requires three things simultaneously: under the 2 GB ceiling, at least three times smaller than the other side, and preferSortMergeJoin flipped to false. That conjunction, with the last condition off by default, is why it is rare in practice.

Which build side each join type allows

This is the rule that makes broadcast joins look unpredictable, and it is pure lookup, not heuristics.

Join type Broadcast the left? Broadcast the right?
Inner, cross Yes Yes
Left outer No Yes
Right outer Yes No
Full outer No No
Left semi, left anti No Yes

The logic follows from correctness. A left outer join must emit every left row, so the left side has to be streamed and the right one built. A right outer is the mirror image. A full outer must emit unmatched rows from both sides, so neither can be the build side, and no full outer join is ever a broadcast hash join regardless of how small either side is.

Shuffle hash join is less restricted, which is why the strategy table says it supports all join types: its build side may be the left for inner, left outer, full outer and right outer, and the right for those plus left semi, left anti and existence joins.

How do you make each strategy appear?

Every rule above is testable. This section builds two tables, forces all five strategies, and reads back the one that actually executed. The transcripts below are from Spark 4.1.3 with spark.sql.shuffle.partitions set to 8, and the five selection predicates quoted earlier were checked against the 4.1.3 source.

A fact table and a dimension:

from pyspark.sql import SparkSession, functions as F

spark = (SparkSession.builder
         .appName("join-strategies")
         .master("local[4]")
         .config("spark.sql.shuffle.partitions", "8")
         .getOrCreate())

cities = (spark.range(0, 600)
          .select(F.col("id").alias("city_id"),
                  F.concat(F.lit("city-"), F.col("id")).alias("city_name"),
                  F.when(F.col("id") % 3 == 0, "IN")
                   .when(F.col("id") % 3 == 1, "US")
                   .otherwise("DE").alias("country_code")))

trips = (spark.range(0, 3_000_000)
         .select(F.col("id").alias("trip_id"),
                 (F.col("id") % 500).alias("city_id"),
                 (F.col("id") % 90000).alias("rider_id"),
                 (F.col("id") % 4000 / 100.0).cast("decimal(10,2)").alias("fare_amount")))

Two helpers. The first reports the strategy that ran, the second reports the size Catalyst believes a side to be, which is the input to every size rule:

def strategy(df):
    plan = df._jdf.queryExecution().executedPlan().toString()
    for name in ["BroadcastHashJoin", "ShuffledHashJoin", "SortMergeJoin",
                 "BroadcastNestedLoopJoin", "CartesianProduct"]:
        if name in plan:
            return name
    return "unknown"

def size_in_bytes(df):
    return int(df._jdf.queryExecution().optimizedPlan().stats().sizeInBytes())

For these two tables the estimates are:

size_in_bytes(cities) =     16,800 bytes
size_in_bytes(trips)  = 60,000,000 bytes

trips references only the first 500 city_id values, so 100 of the 600 cities have no trips. That matters for the outer joins later. cities is far under the 10 MB broadcast threshold and trips is far over it, so the defaults should broadcast. Now force each strategy in turn:

# 1 broadcast hash join, the default here because cities is under the threshold
strategy(trips.join(cities, "city_id"))

# 2 sort-merge, by taking broadcasting away
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")
strategy(trips.join(cities, "city_id"))
spark.conf.unset("spark.sql.autoBroadcastJoinThreshold")

# 3 shuffle hash, by asking for it
strategy(trips.join(cities.hint("SHUFFLE_HASH"), "city_id"))

# 4 broadcast nested loop, by using an inequality instead of an equality
strategy(trips.join(cities, trips.city_id > cities.city_id))

# 5 cartesian, a deliberate cross product with broadcasting off
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")
strategy(trips.limit(1000).crossJoin(cities.limit(50)))
spark.conf.unset("spark.sql.autoBroadcastJoinThreshold")

The operator line each one produced:

1  BroadcastHashJoin [city_id#6L], [city_id#1L], Inner, BuildRight, false
2  SortMergeJoin [city_id#6L], [city_id#1L], Inner
3  ShuffledHashJoin [city_id#6L], [city_id#1L], Inner, BuildRight
4  BroadcastNestedLoopJoin BuildRight, Inner, (city_id#6L > city_id#1L)
5  CartesianProduct

All five, from one dataset, by changing one thing at a time. Note that case 4 returned 748,500,000 rows from a 3,000,000-row table: an inequality against 600 cities is a near-cross-product, which is what the strategy name is telling you.

Each of the four hints behaves as the procedure says, overriding the size rules but not the join-type rules:

Hint Written as Strategy it produced
BROADCAST cities.hint("BROADCAST") BroadcastHashJoin
MERGE cities.hint("MERGE") SortMergeJoin
SHUFFLE_HASH cities.hint("SHUFFLE_HASH") ShuffledHashJoin
SHUFFLE_REPLICATE_NL cities.hint("SHUFFLE_REPLICATE_NL") CartesianProduct

A hint is the one lever that skips the size checks entirely, which is why BROADCAST on a large table moves the failure to the driver rather than preventing it. The last row is the most destructive of the four: it printed CartesianProduct (city_id#6L = city_id#1L), meaning the equality was demoted from a join key to a filter applied after the cross product was built.

The join type overrides the size

The size rules are only consulted after the join type has had its say, and this is where the surprises live. One rule sits underneath all of them: an outer join must emit every row of its preserved side, so that side has to be streamed, and the build side is always the other one.

A left outer preserves the left and therefore builds the right. A right outer preserves the right and builds the left. A full outer preserves both, so it can build neither. That is the entire build-side table, derived rather than memorised.

Four joins over the same two tables make the consequence concrete:

strategy(trips.join(cities, "city_id", "right"))   # preserves cities, the small side
strategy(cities.join(trips, "city_id", "left"))    # preserves cities as well
strategy(cities.join(trips, "city_id", "right"))   # preserves trips, the large side
strategy(trips.join(cities, "city_id", "left"))    # preserves trips as well
trips  RIGHT JOIN cities  3,000,100 rows  SortMergeJoin     [city_id#6L], [city_id#1L], RightOuter
cities LEFT  JOIN trips   3,000,100 rows  SortMergeJoin     [city_id#1L], [city_id#6L], LeftOuter
cities RIGHT JOIN trips   3,000,000 rows  BroadcastHashJoin [city_id#1L], [city_id#6L], RightOuter, BuildLeft, false
trips  LEFT  JOIN cities  3,000,000 rows  BroadcastHashJoin [city_id#6L], [city_id#1L], LeftOuter, BuildRight, false

The first two lines are one query written two ways. Both preserve cities, both return 3,000,100 rows, and comparing them row by row finds no difference in either direction. Both sort-merge. The last two are also one query written two ways, both preserve trips, and both broadcast.

The two pairs, though, are not the same query. They differ by exactly the 100 cities that have no trips. That is the trap in the advice usually given here: swapping the argument order of an outer join changes which side is preserved, so it changes the answer. It is a join-type change wearing the costume of a refactor, and it is the one rewrite that can quietly drop rows.

So the honest rule is narrower than “swap the sides”:

  • If you must preserve the large side, the build side is the small one, and the broadcast comes for free.
  • If you must preserve the small side, the build side is the large one, and no threshold, hint or reordering produces a broadcast hash join. Sort-merge is the correct plan, and the lever you have is reducing the large side, not persuading the optimiser.

Two more results from the same tables close the picture:

trips FULL JOIN cities           3,000,100 rows  SortMergeJoin     ... FullOuter
trips LEFT ANTI JOIN cities              0 rows  BroadcastHashJoin ... LeftAnti, BuildRight, false

A full outer preserves both sides, so it sort-merges however small cities is. A left anti preserves the left and returns no columns from the right at all, so the right is always available as the build side and it broadcasts. It returns zero rows here because every trip does have a matching city, which is the correct answer to “trips whose city is missing”.

Why shuffle hash join is rare, demonstrated

The three conditions are easier to believe when you watch them fail one at a time. Take a 400,000-row riders table against trips on rider_id:

riders = (spark.range(0, 400_000)
          .select(F.col("id").alias("rider_id"),
                  F.concat(F.lit("rider-"), F.col("id")).alias("rider_name")))

def attempt(threshold, prefer_smj):
    spark.conf.set("spark.sql.autoBroadcastJoinThreshold", str(threshold))
    spark.conf.set("spark.sql.join.preferSortMergeJoin", prefer_smj)
    return strategy(trips.select("trip_id", "rider_id").join(riders, "rider_id"))

attempt(-1,          "false")   # the usual "disable broadcast" trick
attempt(2 * 1024**2, "false")   # broadcast still off, but a positive threshold
attempt(2 * 1024**2, "true")    # preferSortMergeJoin back at its default

With riders at 7,200,000 bytes and the projected trips at 36,000,000:

Attempt canBroadcastBySize Hash-map ceiling canBuildLocalHashMapBySize muchSmaller preferSortMergeJoin Strategy
threshold = -1 false -1 x 8 = -8 false true false SortMergeJoin
threshold = 2 MB false 2 MB x 8 = 16 MB true true false ShuffledHashJoin
threshold = 2 MB, prefer on false 2 MB x 8 = 16 MB true true true SortMergeJoin

The middle row is a shuffle hash join with no hint at all, which is the thing most people have never seen. The other two rows each break one condition.

The first row deserves its own warning, because it is a trap in the most common tuning gesture there is. Setting autoBroadcastJoinThreshold to -1 does not just disable broadcasting, it also makes shuffle hash join unreachable, since canBuildLocalHashMapBySize multiplies that same threshold by the shuffle partition count and -1 x 8 is negative. Every size below zero fails the test. If you want to stop broadcasting but keep shuffle hash join available, set the threshold to a small positive number instead of -1.

What does adaptive execution change about joins?

Everything above is plan-time. Adaptive query execution re-plans after a shuffle stage completes, when the sizes are measurements rather than estimates, and spark.sql.adaptive.enabled defaults to true.

For joins it does two things:

  1. Demotes a sort-merge join to a broadcast hash join when a side turns out small enough, using spark.sql.adaptive.autoBroadcastJoinThreshold.
  2. Splits skewed partitions, with spark.sql.adaptive.skewJoin.enabled default true.

The skew thresholds are strict and worth knowing, because they explain “skew handling is on but nothing happened”. A partition is treated as skewed only when it is larger than spark.sql.adaptive.skewJoin.skewedPartitionFactor, default 5.0, times the median and larger than spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes, default 256MB. On a job whose partitions are all under 256 MB, it never fires no matter how uneven they are.

The practical consequence for reading plans: explain() shows the planned shape, not the final one. Before execution the tree is wrapped in AdaptiveSparkPlan with isFinalPlan=false. The SQL tab of the UI after the query finishes shows what actually ran.

How do you read the join out of a plan?

df.explain("formatted")

Here is the real output for the broadcast join above, projected down to three columns. The tree comes first, then a numbered detail block per operator:

== Physical Plan ==
AdaptiveSparkPlan (9)
+- Project (8)
   +- BroadcastHashJoin Inner BuildRight (7)
      :- Project (3)
      :  +- Filter (2)
      :     +- Range (1)
      +- BroadcastExchange (6)
         +- Project (5)
            +- Range (4)

(6) BroadcastExchange
Input [2]: [city_id#1L, city_name#2]
Arguments: HashedRelationBroadcastMode(List(input[0, bigint, false]),false), [plan_id=28]

(7) BroadcastHashJoin
Left keys [1]: [city_id#6L]
Right keys [1]: [city_id#1L]
Join type: Inner
Join condition: None

(9) AdaptiveSparkPlan
Output [3]: [trip_id#5L, city_name#2, fare_amount#8]
Arguments: isFinalPlan=false

Four things in that output are worth reading deliberately. BuildRight on node 7 names the build side. BroadcastExchange on node 6 appears under the right branch only, so the left branch has no exchange and the large side is not shuffled. Node 6 lists city_id and city_name but not country_code, because the projection pushed down and shrank what gets broadcast. And node 9 says isFinalPlan=false, so this is the plan before adaptive execution has had a say.

What to look for, in order of how much it tells you:

In the plan Meaning
BroadcastHashJoin The small side is being broadcast. Usually what you want
SortMergeJoin Both sides shuffled and sorted. Correct for large-to-large
ShuffledHashJoin Rare without a hint
BroadcastNestedLoopJoin Non-equi, or a missing condition
CartesianProduct Almost always a mistake
Exchange hashpartitioning(k, 200) The shuffle, with its key and partition count
BroadcastExchange The broadcast being built
ReusedExchange One shuffle serving two branches. Good news
isFinalPlan=false AQE has not re-planned yet

A useful habit: count the Exchange nodes. That is the number of shuffles, and if two consecutive operations key on the same column, one exchange can serve both.

Which join strategy fits which scenario?

Scenario Aim for How
Large fact, small dimension Broadcast hash Nothing, if the dimension is under 10 MB. Filter and project it first if not
Large fact, medium dimension (tens of MB) Broadcast hash Reduce it first, then raise the threshold deliberately if it is genuinely small in memory
Large to large, recurring Sort-merge, or no shuffle at all Bucket both tables on the key with the same bucket count
Large to large, one-off Sort-merge Filter and project both sides, size the shuffle, handle skew
Full outer with a small side Sort-merge, unavoidable Accept it, or restructure as two outer joins and a union
Outer join preserving the large side Broadcast hash Nothing. The build side is the small one already
Outer join preserving the small side Sort-merge, unavoidable The preserved side cannot be built. Reduce the other side instead
Range or inequality condition Broadcast nested loop Add an equality that buckets both sides, so the nested loop runs within buckets
Existence check Semi or anti join left_semi or left_anti, never inner plus distinct
Skewed join key Sort-merge plus skew handling AQE first, then salt the hot key

Two of those need expanding.

The range join. A condition like ON t.ts BETWEEN d.start AND d.end has no equality, so Spark can only nested-loop it. The standard fix is to manufacture a key: bucket both sides by a coarser value, join on that equality and the range, so the comparison happens within each bucket rather than across the whole product.

(events.withColumn("day", F.to_date("ts"))
   .join(windows.withColumn("day", F.to_date("start")),
         ["day"])                                    # an equality to partition on
   .filter((F.col("ts") >= F.col("start")) & (F.col("ts") <= F.col("end"))))

The full outer with a small side. If broadcasting matters more than the single scan, a full outer can be rewritten as a left outer union an anti join the other way, both of which can broadcast. It is more code and two passes, so only worth it when the sort-merge is genuinely the bottleneck.

Common misconceptions

“Raising autoBroadcastJoinThreshold will make this broadcast.” Not if the join type forbids it. A full outer join can never be a broadcast hash join, because both sides must emit unmatched rows and neither may be the build side. No threshold changes that. Check the join type before the size.

“Swapping the sides of an outer join is a free optimisation.” It changes the answer. An outer join preserves one side, and swapping the argument order swaps which side that is. Measured on 600 cities where 100 have no trips, the two orderings differ by exactly those 100 rows, while returning the same count on data where everything matches, which is how the bug hides.

“Setting autoBroadcastJoinThreshold to -1 only turns off broadcasting.” It also removes shuffle hash join from consideration. The shuffle-hash ceiling is autoBroadcastJoinThreshold * spark.sql.shuffle.partitions, so -1 makes it negative and every candidate fails the size test. Use a small positive value if you want to stop broadcasting and keep the option.

“A different join strategy can change my results.” It cannot. Strategies are algorithms for the same semantics. If the rows changed, the join type or the condition changed, or nulls in the key are behaving as SQL says they should.

Production tips

  • Read the plan before tuning. The strategy Spark chose tells you which rule fired, and therefore which lever exists.
  • Reduce before joining. Filtering and projecting a side is the highest-value change available, because it can move the join into the broadcast band and change the strategy entirely.
  • Do not hint BROADCAST on something large. A hint overrides the size check, so you move the failure to the driver, which must collect it first.
  • Prefer a small positive autoBroadcastJoinThreshold over -1. Setting it to -1 also drives the shuffle-hash ceiling negative, removing that strategy from consideration along with broadcasting.
  • Treat a full outer join as unbroadcastable. Raising the threshold will never help; the join type is what forbids it.
  • For an outer join, ask which side is preserved, not which side is small. The build side is always the non-preserved one, so an outer join that must preserve the small side cannot broadcast at any threshold. Reordering the arguments changes the result, not just the plan.
  • Fix the estimate rather than forcing the plan where you can, with ANALYZE TABLE ... COMPUTE STATISTICS, so Catalyst chooses correctly on its own.
  • Bucket the keys you join on repeatedly. It is the only way to remove the shuffle from a large-to-large join permanently.
  • Trust the SQL tab over explain() for what actually ran, since AQE re-plans after the shuffle.

Frequently asked questions

Why is my small table not being broadcast? Four common reasons: the estimate is of the unfiltered table because the filter did not push down; the table is compressed on disk and larger in memory than you think; there are no statistics so Spark fell back to file size; or the join type forbids that side as the build side. Check the last one first, because no config change fixes it.

Does a bigger autoBroadcastJoinThreshold make joins faster? Sometimes, and it moves the failure mode from slow to fragile. Every executor holds a full copy and the driver assembles it first, so a 500 MB broadcast across 50 executors is 25 GB of cluster memory plus a driver that must hold it. Raise it deliberately, with a number you can justify.

Why do I almost never see ShuffledHashJoin? Because it needs three conditions at once and one of them, spark.sql.join.preferSortMergeJoin being false, is not the default. Spark prefers sort-merge because it cannot run out of memory, while a hash join can. There is also a trap: if you disabled broadcasting with autoBroadcastJoinThreshold = -1, the shuffle-hash ceiling became negative too, so the strategy cannot be chosen without a hint no matter what else you set. The table above shows the same join picking it once the threshold is a small positive number.

Can a join strategy change my results? No. Strategies are algorithms for the same semantics. If results changed, the join type or the condition changed, or there are nulls in the key behaving as SQL says they should.

Why is my join a CartesianProduct when I wrote a condition? Because Catalyst could not use the condition as an equality: it is an inequality, it is wrapped in an expression it cannot match on, or it sits in a where after a cross join in a way that could not be pushed into the join.

What about joining on a collated string column? In Spark 4 this matters. Hash joins require keys that are binary-stable, and a case- or accent-insensitive collation is not, because two rows can be equal without having equal bytes. Spark does not give up on hashing for that, though: it rewrites each key into collationkey(k), a binary-stable form, and hashes that. On Spark 4.1.3 a join on two UTF8_LCASE columns broadcast by default, took ShuffledHashJoin under a SHUFFLE_HASH hint, and sort-merged only once broadcasting was off, with both join keys printed as collationkey(k) in every plan. Hash joins stay available; the price is a collationkey evaluated on every row of both sides.

Conclusion

The thing worth taking away is that join selection is a lookup, not a judgement. Spark asks, in a fixed order: is there a hint, is this an equi-join, is one side small enough, does the join type permit that side as the build side, are the keys sortable. The first rule that fits wins. Nothing about it is adaptive except the parts AQE re-plans after a shuffle.

That reframes tuning. You are not persuading an optimiser, you are changing the inputs to a decision procedure. Filtering a dimension until it fits under 10 MB changes the answer to “is one side small enough”. Deciding which side an outer join preserves settles which side may be built, though it changes the result too, so it is a modelling decision rather than a tuning knob. Bucketing changes whether a shuffle is needed at all. Each of those is a specific rule with a specific lever.

And three answers are simply no. A full outer join will not broadcast, a non-equi join will not hash, and an outer join that must preserve its small side will not broadcast either. When you hit those, the fix is to change the query shape, or to accept the plan and shrink the other side, rather than to keep raising thresholds.

References

Trademarks

Apache Spark, Apache Hive, Apache Parquet and Apache are either registered trademarks or trademarks of The Apache Software Foundation in the United States and other countries.

Found this useful?

These posts and tools are free. If one saved you an afternoon, you can buy me a coffee.

Buy me a coffee