On a cluster, matching rows must first meet on one executor. The planner picks a strategy from the join condition, the estimated sizes and any hints.
| Strategy | Moves | Needs | Best for |
|---|---|---|---|
| Broadcast hash join | Small side to every executor | Equality; one side under 10 MB | Fact to dimension |
| Shuffle hash join | Both sides, by key | Equality; one side fits per partition | Medium build side |
| Sort-merge join | Both sides, by key, then sorts | Equality; sortable keys | Two large tables (default) |
| Broadcast nested loop | Small side everywhere | Any condition | Small non-equi joins |
import contextlib, io
def strategy(sql):
buf = io.StringIO()
with contextlib.redirect_stdout(buf):
spark.sql(sql).explain() # the plan AQE starts from
return next(w for w in buf.getvalue().split() if w.endswith("Join"))
EQUI = """SELECT {} c.country, sum(l.qty) FROM lines l
JOIN customers c ON l.customer_id = c.customer_id GROUP BY c.country"""
for hint in ["", "/*+ MERGE(c) */", "/*+ SHUFFLE_HASH(c) */", "/*+ BROADCAST(c) */"]:
print(f"{hint or 'no hint':<24} {strategy(EQUI.format(hint))}")
print(f"{'price range, no equality':<24} "
f"{strategy('SELECT * FROM lines l JOIN books b ON l.unit_price > b.price')}")Output
no hint BroadcastHashJoin /*+ MERGE(c) */ SortMergeJoin /*+ SHUFFLE_HASH(c) */ ShuffledHashJoin /*+ BROADCAST(c) */ BroadcastHashJoin price range, no equality BroadcastNestedLoopJoin
The pruned customer table came in under the 10 MB spark.sql.autoBroadcastJoinThreshold, so it was broadcast unasked. Without an equality condition only nested loops remain, comparing every pair: fine against six books, ruinous between two large tables. Conflicting hints rank BROADCAST over MERGE over SHUFFLE_HASH.