Every join in Spark SQL becomes one of a handful of physical algorithms. Which one you get decides whether a query shuffles two terabytes or two megabytes, whether it runs for four minutes or forty, and whether an executor dies of an out-of-memory error at 3 a.m. Most of the time the planner chooses well from table sizes and configuration. Join strategy hints exist for the rest of the time: when statistics are missing or wrong, when you know something about the data the optimizer cannot, or when a plan needs to be pinned so it stops changing between runs.
This article explains the strategies, how the planner chooses among them, hint precedence, the cases where Spark ignores a hint, and how hints interact with Adaptive Query Execution, then works through a three-table join and ends with a checklist. For the internals of the broadcast path itself, see Spark broadcast joins.
The four physical join strategies
A join has two inputs, a join type (inner, left outer, right outer, full outer, left semi, left anti, cross) and a condition. When the condition contains equality predicates on columns from both sides, Spark calls them equi-join keys, and three hash- or sort-based algorithms become available. Without equi-join keys, only nested-loop algorithms remain.
Broadcast hash join. One side, the build side, is collected to the driver, turned into an in-memory hash table and broadcast to every executor. Each task then streams its partition of the other side through the table. The large side is never shuffled, which is why this is usually the fastest option, but the whole build side must fit in memory on the driver and on every executor.
Shuffled hash join. Both sides are shuffled by the join keys so matching rows land in the same partition. Each task builds a hash table from its slice of the build side and probes it with its slice of the other side. There is no sort, but each task's build slice must fit in its memory.
Sort-merge join. Both sides are shuffled by key and then sorted within each partition; the task walks both sorted streams in step. Sorting costs CPU, but it can spill to disk, so it scales to inputs of any size. It is Spark's default for large equi-joins.
Nested-loop joins. Broadcast nested loop and Cartesian product (inner joins only) compare every pair of rows against an arbitrary condition. Both are quadratic; they are what non-equi joins get.
| Strategy | Shuffles | Memory risk | Needs equi-keys |
|---|---|---|---|
| Broadcast hash | Neither side | Whole build side on driver and executors | Yes |
| Shuffled hash | Both sides | Build slice per task | Yes |
| Sort-merge | Both sides | Low, sort spills | Yes, sortable |
| Broadcast nested loop | Neither side | Whole broadcast side | No |
| Cartesian product | Both sides | Low, but quadratic work | No, inner only |
How the planner chooses, with and without hints
Without hints, the planner's rules for an equi-join run roughly in this order. If one side's estimated size is below spark.sql.autoBroadcastJoinThreshold (10 MB by default; -1 disables automatic broadcasting) and the join type allows that side to be the build side, it plans a broadcast hash join. Otherwise, if spark.sql.join.preferSortMergeJoin is false and one side is small enough to build a per-partition hash map and much smaller than the other, it plans a shuffled hash join. Otherwise, if the keys are sortable, it plans a sort-merge join. Non-equi joins fall through to Cartesian product (inner joins) or broadcast nested loop.
Hints are consulted before those size rules. The pseudocode below captures the shape of the planner's JoinSelection strategy, which has more cases.
def select_join(join):
hinted = strategy_hints(join) # from both sides, highest priority first
for strategy in ["BROADCAST", "MERGE", "SHUFFLE_HASH", "SHUFFLE_REPLICATE_NL"]:
if strategy in hinted and supports(strategy, join.type, join.has_equi_keys):
return plan(strategy, build_side=pick_build_side(join, hinted))
if hinted:
log_warning("hint not supported for this join; ignoring")
# no usable hint: size and config rules
if join.has_equi_keys:
if can_broadcast_by_size(join): # estimate < autoBroadcastJoinThreshold
return plan("BROADCAST")
if not prefer_sort_merge and small_enough_for_local_hash_map(join):
return plan("SHUFFLE_HASH")
if keys_sortable(join):
return plan("MERGE")
if join.type == "inner":
return plan("CARTESIAN")
return plan("BROADCAST_NESTED_LOOP")Two consequences follow. A hint overrides size: BROADCAST on a 500 MB table broadcasts it at fifty times the threshold. A hint never overrides correctness: if the strategy cannot execute that join type, the hint is dropped.
Writing hints in SQL and the DataFrame API
In SQL a hint is a comment immediately after SELECT and names the relation, by table name or alias, that it applies to. The Spark documentation lists four join hints and their aliases: BROADCAST (also BROADCASTJOIN and MAPJOIN), MERGE (also SHUFFLE_MERGE and MERGEJOIN), SHUFFLE_HASH and SHUFFLE_REPLICATE_NL. Separate partitioning hints, COALESCE, REPARTITION, REPARTITION_BY_RANGE and REBALANCE, shape output partitions rather than join algorithms.
-- SQL: hint the alias, not the underlying table name, when an alias is used
SELECT /*+ BROADCAST(c) */ o.order_id, c.segment
FROM orders o
JOIN countries c ON o.country_code = c.code;
-- several hints in one comment
SELECT /*+ SHUFFLE_HASH(p), MERGE(o) */ *
FROM orders o JOIN payments p ON o.order_id = p.order_id;# PySpark: three equivalent ways to ask for a broadcast of `countries`
from pyspark.sql.functions import broadcast
joined = orders.join(broadcast(countries), orders.country_code == countries.code)
joined = orders.join(countries.hint("broadcast"), orders.country_code == countries.code)
joined = orders.join(countries.hint("BROADCASTJOIN"), orders.country_code == countries.code)
# other strategies use .hint() with the strategy name
joined = orders.join(payments.hint("shuffle_hash"), "order_id")
joined = orders.join(payments.hint("merge"), "order_id")
joined.explain(mode="formatted") # look for BroadcastHashJoin / ShuffledHashJoin / SortMergeJoinHint names are case-insensitive. A SQL hint must name the relation as that query block references it; a hint that finds nothing to attach to only produces a warning, so always confirm the result with EXPLAIN.
Precedence, build sides and the cases Spark ignores
When different strategy hints appear on the two sides of the same join, the documentation states the priority: BROADCAST over MERGE over SHUFFLE_HASH over SHUFFLE_REPLICATE_NL. So /*+ MERGE(a), BROADCAST(b) */ produces a broadcast of b if the join type allows it. When both sides carry BROADCAST or both carry SHUFFLE_HASH, Spark chooses the build side from the join type and the estimated sizes.
The join type constrains which side can be built. For a broadcast hash join, the build side must be the side whose unmatched rows are not emitted:
| Join type | Broadcast left | Broadcast right | Broadcast hash possible? |
|---|---|---|---|
| Inner | Yes | Yes | Yes, either side |
| Left outer, left semi, left anti | No | Yes | Only by broadcasting the right side |
| Right outer | Yes | No | Only by broadcasting the left side |
| Full outer | No | No | No; shuffled hash (Spark 3.1+) or sort-merge |
This is the most common reason a hint is ignored. SELECT /*+ BROADCAST(o) */ ... FROM orders o LEFT JOIN countries c asks Spark to broadcast the preserved side of a left outer join, which a hash join cannot do: every row of orders must appear in the output even without a match, and a task that sees only part of the probe side cannot know which build rows went unmatched elsewhere. Spark logs that the hint is not supported for that join type and plans something else.
Other cases: MERGE needs equi-join keys of sortable types; SHUFFLE_REPLICATE_NL works only for inner joins; and Spark refuses to broadcast a relation larger than 8 GB however it is hinted. A hint is a strong preference applied within the rules, never a command that bypasses them.
Hints and Adaptive Query Execution
With spark.sql.adaptive.enabled (on by default since Spark 3.2), Spark plans a query, executes it stage by stage and re-optimises the remaining plan using the actual sizes of completed shuffle stages. Two conversions affect joins. AQE converts a sort-merge join to a broadcast hash join when the runtime size of one side falls below spark.sql.adaptive.autoBroadcastJoinThreshold, which defaults to the same value as the static threshold. It can also prefer a shuffled hash join when every partition is below spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold, which defaults to 0, so that conversion is off unless you set it.
AQE also splits skewed partitions in sort-merge joins. A partition counts as skewed when it is larger than both skewedPartitionFactor (5.0) times the median partition and skewedPartitionThresholdInBytes (256 MB). This often removes the need for manual salting; the AQE deep dive covers the mechanics.
The practical rule is that AQE fixes estimates that were wrong because a filter or aggregate shrank data more than expected, which is exactly the situation people used to fix with broadcast hints. Before adding a hint, check whether AQE already converts the join at runtime: run the query, open the SQL tab of the Spark UI, and compare the initial plan with the final plan, which shows isFinalPlan=true. Strategy hints remain attached to the logical join, so in current releases AQE re-planning honours them; confirm that on your version by inspecting the final plan, because a hint written for Spark 2 can pin a worse plan than AQE would choose.
Worked example: a three-table join, step by step
Consider a nightly job joining orders (2 TB of Parquet, about 8 billion rows), customers (60 GB) and countries (250 rows, 40 KB). The query keeps orders from the last day, roughly 25 GB, and filters customers to one market, roughly 300 MB on disk.
SELECT o.order_id, o.amount, cu.tier, co.region
FROM orders o
JOIN customers cu ON o.customer_id = cu.customer_id
JOIN countries co ON cu.country_code = co.code
WHERE o.order_date = DATE'2026-09-30'
AND cu.market = 'EU';Default plan. countries is far below 10 MB, so it is broadcast automatically. The static estimate for filtered customers is unreliable: without column statistics the planner cannot know how selective market = 'EU' is, and it plans a sort-merge join between 25 GB of orders and the customers side. Both are shuffled and sorted.
With AQE. AQE measures about 300 MB of customers shuffle output, above the 10 MB adaptive threshold, so the sort-merge join stays; AQE only splits any skewed partitions.
With a hint. 300 MB of compressed Parquet might expand to 1 to 2 GB as a Java hash table, depending on column types. On executors with 16 GB of heap that is affordable, so /*+ BROADCAST(cu) */ removes the 25 GB shuffle and both sorts. It is only safe because the filter is stable: if market = 'EU' is ever replaced with a parameter that selects every market, the same hint broadcasts 60 GB and fails with a broadcast timeout or out-of-memory error.
A safer alternative. Raise spark.sql.adaptive.autoBroadcastJoinThreshold for this job to 512 MB and leave the static threshold alone. The decision then uses measured size, and the broadcast happens only when the data really is small. Use the hint only if the plan must be pinned, and document the size assumption next to it.
Failure modes and how to diagnose them
- Driver out of memory during broadcast. The build side is collected to the driver first, so a hinted table larger than expected, or much larger in memory than on disk, takes the driver down, sometimes as a
spark.driver.maxResultSizefailure. Fix the size assumption, not the driver memory. - Broadcast timeout.
spark.sql.broadcastTimeout(300 seconds by default) also fires when a small build side comes from a slow upstream computation. Materialise the build side or raise the timeout for that job. - Shuffled hash join OOM. One skewed key can make one partition's build slice huge, and the task fails where a sort-merge join would have spilled. Use SHUFFLE_HASH only when per-partition sizes are known and even.
- Hint silently dropped. Wrong alias, unsupported join type or a view hiding the relation. Check driver-log warnings and
EXPLAIN. - Stale hints. A dimension hinted for broadcast years ago has grown tenfold. Record each hint's size assumption and alert on it.
- Accidental nested loops. A condition with an
ORacross keys, or only a range predicate, has no equi-keys, so a nested-loop join appears. Rewrite with an equality where possible.
A plan assertion in tests catches a pinned strategy that regressed after an upgrade or schema change.
def assert_join_strategy(df, expected):
"""expected: 'BroadcastHashJoin', 'SortMergeJoin', 'ShuffledHashJoin', ..."""
plan = df._jdf.queryExecution().executedPlan().toString()
found = [s for s in ("BroadcastHashJoin", "ShuffledHashJoin", "SortMergeJoin",
"BroadcastNestedLoopJoin", "CartesianProduct") if s in plan]
if expected not in found:
raise AssertionError(f"expected {expected}, plan has {found}")With AQE on, call an action before asserting if you need the final strategy. _jdf is an internal handle, so keep this in tests.
Trade-offs: hints, statistics or configuration
Three levers improve a join, and they age differently. Statistics (ANALYZE TABLE ... COMPUTE STATISTICS FOR COLUMNS with spark.sql.cbo.enabled) help every query on the table but must be refreshed. Configuration, such as a per-job adaptive broadcast threshold, adapts to measured sizes and is the right default. Hints are precise and immediate, but are frozen assumptions in query text.
Use hints for knowledge the engine cannot measure: a stable, known filter selectivity, a dimension small by design, or a plan that must not change. Avoid them in shared views and library code, where one caller's size assumption is imposed on all. Runtime techniques such as bloom-filter join pruning complement hints by shrinking the large side before the shuffle, and the shuffle architecture explains what a sort-merge join actually pays for. For how hints fit into the wider rule pipeline, see the Spark SQL optimizer.
What to do next
- Pick your slowest job and run
EXPLAIN FORMATTED; note every join node and which side is built or broadcast. - Run it once with AQE on and compare the initial and final plans in the Spark UI SQL tab before touching any hint.
- For joins that stay sort-merge although one side is small after filtering, try raising
spark.sql.adaptive.autoBroadcastJoinThresholdfor that job first. - Hint only where the size assumption is stable, with a comment stating the expected size.
- Never hint the preserved side of an outer join; check the build-side table above.
- Add a plan assertion test for each pinned join and rerun it on every Spark upgrade.
- Audit hints quarterly against table sizes; delete those AQE handles.