When Spark joins two large tables on an equality condition, the default physical plan is a sort-merge join. Both sides are shuffled so that rows with the same key land in the same partition, each partition is sorted by the key, and the two sorted streams are merged in one pass. It is the join that works when neither side fits in memory, which is why it is the default, and it is also where many slow stages come from: the shuffle, the sort and, on skewed or duplicated keys, a buffer that spills.
How Spark chooses between join strategies and how hints change that choice are covered in Spark join hints. This page covers what happens once a sort-merge join is chosen: the merge algorithm, which side Spark buffers and why, how to read it in a plan, how to remove its shuffle and sort with bucketing, and how to diagnose and fix the cases where it goes wrong. Configuration defaults are from the Spark 3.5 source. Plan fragments are abridged illustrations of the shape you will see, not captured output.
The plan in one picture
The merge algorithm
The merge itself is the classic algorithm from databases. With both inputs sorted on the key, keep a cursor on each. If the left key is smaller, advance the left; if the right key is smaller, advance the right; if they are equal, every left row with that key pairs with every right row with that key. The only state is the group of rows on one side that share the current key, and it is needed because the other side may have several rows with that key too.
def sort_merge_join(streamed, buffered): # both iterators sorted by key
s, b = next(streamed, None), next(buffered, None)
while s is not None and b is not None:
if s.key is None or s.key < b.key: # null keys never match an equi-join
s = next(streamed, None)
elif s.key > b.key:
b = next(buffered, None)
else:
k, group = s.key, []
while b is not None and b.key == k: # collect every buffered row for k
group.append(b)
b = next(buffered, None)
while s is not None and s.key == k: # each streamed row meets the group
for g in group:
yield s, g
s = next(streamed, None)Each input row is read once and the memory needed is one key group, not a whole table, which is why the join scales to inputs far larger than memory. The price is that both sides must be partitioned and sorted on the key first, and that a single key with very many buffered rows makes the group large.
Inside SortMergeJoinExec
In Spark the operator is SortMergeJoinExec. It requires each child to be clustered by its join keys, which the planner satisfies by inserting an Exchange hashpartitioning on each side with the same partition count, and sorted ascending by the keys, which adds a Sort on each side. Partition i of the left then contains exactly the keys of partition i of the right, and one task joins each pair.
One side is streamed and the other buffered. For inner, left outer, full outer, left semi and left anti joins the left child is streamed and the right is buffered; for a right outer join the roles swap. The buffered rows for the current key go into an ExternalAppendOnlyUnsafeRowArray. Rows are copied into an in-memory array until their count reaches spark.sql.sortMergeJoinExec.buffer.in.memory.threshold; only then does the buffer switch to a sorter that can spill to disk, governed by spark.sql.sortMergeJoinExec.buffer.spill.threshold. Both are internal settings, and in the Spark 3.5 source the in-memory default is the maximum array length, so with defaults a large key group effectively stays on the heap. The join node does have a spill size metric, but expect it to read zero unless someone lowered the threshold. The realistic symptom of a huge key group is one task with long GC time or an executor out-of-memory error, not a spill.
Streamed rows whose key contains a null are skipped for inner joins, because null never equals anything; outer joins still emit them with nulls on the other side. The two sorted children are executed as separate inputs to the join; whole-stage code generation then fuses the merge loop with the operators above it, such as a projection or a partial aggregate, which is covered in whole-stage codegen.
Join types, extra conditions and duplicate keys
The merge loop is the same for every join type; what changes is what happens to a row without a partner. For a left outer join, a streamed row whose key has no buffered group is emitted once with nulls for the right columns. A left semi join emits the streamed row once if a group exists and never looks at the group's contents beyond that, so duplicates on the buffered side do not multiply output. A left anti join emits the streamed row only when no group exists. A full outer join has to emit unmatched rows from both sides, so Spark uses a separate scanner for it that advances both cursors and tracks which buffered rows found a match.
Extra join conditions that are not equalities, such as a date range, are evaluated on each candidate pair inside the key group. They do not reduce the group size, so a join on customer_id plus a date range still buffers every row for that customer. If the range is selective, add a coarse equality on a derived column, such as the month, so that the key itself narrows the group.
Duplicate keys on both sides multiply. If a key has 1,000 rows on the left and 1,000 on the right, the join emits a million rows for that key alone, and every one of them flows into the next operator. This is correct behaviour, not a Spark problem, but it is the most common reason a join stage writes far more output than either input. Check row counts per key on both sides before blaming the engine, and aggregate or deduplicate one side first when the business question only needs one match.
Reading it in a plan
In explain() output or the SQL tab of the UI, a sort-merge join has a recognisable shape. An abridged, illustrative version:
SortMergeJoin [order_id#1], [order_id#20], Inner
:- Sort [order_id#1 ASC NULLS FIRST], false, 0
: +- Exchange hashpartitioning(order_id#1, 200), ENSURE_REQUIREMENTS
: +- Filter isnotnull(order_id#1)
: +- FileScan parquet orders[...]
+- Sort [order_id#20 ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(order_id#20, 200), ENSURE_REQUIREMENTS
+- Filter isnotnull(order_id#20)
+- FileScan parquet payments[...]Read it from the bottom. The optimizer adds isnotnull filters for an inner equi-join, so null keys are dropped before the shuffle. The two Exchange nodes must have the same partition count; 200 is the default of spark.sql.shuffle.partitions, and with adaptive execution the final count is set at run time after coalescing. Under adaptive execution the root shows AdaptiveSparkPlan, and the final plan may replace the join entirely: if a side turns out small enough after the shuffle, it becomes a broadcast hash join, and if spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold is raised from its default of 0 to at least the advisory partition size, and every partition is no larger than it, the join is planned as a shuffled hash join, which skips the sort. A join marked SortMergeJoin(skew=true) has been split by the skew optimization. The explain plans guide covers the rest of the plan vocabulary.
Removing the shuffle and sort with bucketing
If both tables are written bucketed by the join key into the same number of buckets, each bucket already holds exactly one hash partition, and Spark can drop both Exchange nodes. For a table joined every day, paying for that layout once at write time removes the most expensive part of every later join: the shuffle. The Sorts usually remain. In Spark 3.x a bucketed scan does not report its output as sorted, because proving it requires listing files at planning time; that behaviour sits behind the internal flag {c('spark.sql.legacy.bucketedTableScan.outputOrdering')}, which defaults to false. Sorting data within a partition is cheap compared with a shuffle, so this is rarely worth changing.
(orders.write.bucketBy(512, "order_id").sortBy("order_id")
.format("parquet").saveAsTable("sales.orders_b"))
(payments.write.bucketBy(512, "order_id").sortBy("order_id")
.format("parquet").saveAsTable("sales.payments_b"))
j = spark.table("sales.orders_b").join(spark.table("sales.payments_b"), "order_id")
j.explain() # expect no Exchange under SortMergeJoin; Sort nodes may remainThe conditions are strict. Both sides must use the same bucketing hash, the join keys must match the bucket columns, and the counts must be equal, unless spark.sql.bucketing.coalesceBucketsInJoin.enabled is turned on (default false), which lets Spark coalesce the larger count when the ratio is at most maxBucketRatio (default 4). Bucketing is a Spark table property, so readers that ignore it, and table formats with their own partitioning, get no benefit. Plain partitioning by key is covered in Spark partitioning.
Settings that change the join
| Setting (Spark 3.5) | Default | Effect on sort-merge join |
|---|---|---|
spark.sql.join.preferSortMergeJoin | true | Prefer it over shuffled hash join when both apply |
spark.sql.autoBroadcastJoinThreshold | 10MB | Below this estimated size, broadcast replaces it |
spark.sql.shuffle.partitions | 200 | Partition count for both exchanges |
spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold | 0 | If at least the advisory partition size and no partition exceeds it, use shuffled hash join |
spark.sql.adaptive.skewJoin.enabled | true | Split skewed partitions at run time |
spark.sql.bucketing.coalesceBucketsInJoin.enabled | false | Join bucketed tables with different counts |
Worked example: a nightly join with a hot key
A nightly job joins 2 TB of orders with 600 GB of payments on order_id, then aggregates by day. With the default 200 shuffle partitions, each reduce task would receive about 13 GB, which spills heavily during the sort. Targeting about 200 MB per partition gives 2.6 TB divided by 200 MB, roughly 13,000 partitions; with adaptive execution you can set the initial count high and let coalescing reduce it.
After that change the stage runs, but one task takes 40 minutes while the rest take 2, spending most of that time in garbage collection, and on a bad night an executor dies with an out-of-memory error. The stage's task table shows a shuffle read of 30 GB for that partition against a median of 200 MB. The cause is a placeholder order_id used for 6 million test payments, all matching one placeholder order row. The buffered side, payments, has to hold 6 million rows for that one key in memory.
Two fixes apply. Adaptive skew handling, on by default with spark.sql.adaptive.skewJoin.enabled, splits a partition that is larger than skewedPartitionFactor (default 5) times the median and larger than skewedPartitionThresholdInBytes (default 256 MB), and replicates the matching partition of the other side. It splits along map-output blocks, so it can break up even a single hot key when that key's rows were written by many map tasks. For an inner join either side can be split; for a left outer, semi or anti join only the left side; for a right outer join only the right. Splitting still makes the job shuffle and process 6 million useless rows, though, so the real fix here is data: filter the placeholder before the join. When a hot key is legitimate, use salting, described in salting a skewed join.
Failure modes
| Symptom | Cause | Fix |
|---|---|---|
| One task far slower than the rest | Skewed partition or hot key | AQE skew join; remove junk keys; salt |
| One task with long GC time or OOM in the join | Many buffered rows for one key, held in memory | Find the key; filter or salt it |
| Large spill in the Sort | Too few shuffle partitions | Raise partitions or the AQE initial count |
| Output rows far exceed inputs | Duplicate keys on both sides multiply | Deduplicate or aggregate before joining |
| Exchange still present on bucketed tables | Bucket counts or keys differ | Match them, or enable bucket coalescing |
| Executor OOM in many join tasks | Wide rows or too few partitions | Prune columns; raise partitions |
Trade-offs
A broadcast hash join is fastest when one side is small: no shuffle of the large side and no sort. Remedies for skew under any of the three are collected in Spark data skew. A shuffled hash join skips the sort but must build a hash table for each partition in memory, so it fails when a partition is large. Sort-merge join is the slowest of the three when everything fits, and the one that keeps working when nothing does, because its sorts spill and it holds only one key group at a time. That is why spark.sql.join.preferSortMergeJoin defaults to true. Prefer it for large-to-large joins, remove its shuffle with bucketing for repeated joins, and let adaptive execution downgrade it when the data turns out smaller than the statistics said.
What to do next
- Open the SQL tab for your slowest join and confirm it is a SortMergeJoin; note spill on both Sorts and GC time for the slowest tasks.
- Compare max to median shuffle read for the join stage; a ratio above 5 means skew.
- Count rows per join key on both sides for the top keys; remove placeholder and null-like values before the join.
- Set shuffle partitions from data size, about 100 to 200 MB per partition, or rely on adaptive coalescing.
- For tables joined repeatedly on the same key, bucket both by that key with equal counts and check the plan has no Exchange.
- Prune columns before the join so buffered rows are small.