Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Sort-merge join vs shuffle hash join vs broadcast hash join

PySpark · Joins & Data Layout

Sort-merge join vs shuffle hash join vs broadcast hash join

Hardpyspark-63
joinsbroadcast-joinsort-merge-joinshuffle-hash-joinaqe

Question

Which join strategies does Spark have, and how does it choose?

Solution

Spark has a handful of join strategies. It picks one from the join condition, the join type, the estimated sizes of the two sides, and any hints you give. The three you must know are broadcast hash join, sort-merge join and shuffle hash join.

Broadcast hash join

Spark copies the small table to every executor and builds a hash table there. The big table is not shuffled at all, each partition just probes the hash table. It is the fastest option and is chosen when one side is estimated below spark.sql.autoBroadcastJoinThreshold (10 MB by default), or when you hint it:

from pyspark.sql.functions import broadcast
orders.join(broadcast(dim_country), "country_id")

Sort-merge join

The default for large equi-joins. Both sides are shuffled by the join key, sorted within each partition, then merged. It scales to any size because sorting can spill to disk, but it costs two shuffles and two sorts.

Shuffle hash join

Both sides are shuffled, but instead of sorting, Spark builds a hash table from the smaller side in each partition. It can beat sort-merge when one side is much smaller per partition, but it needs the hash table to fit in memory. Spark prefers sort-merge by default for safety, so you usually ask for this one with a hint.

The non-equi cases

If there is no equality condition (a.start < b.ts), hashing is not possible. Spark uses a broadcast nested loop join if one side is small, and a cartesian product otherwise, which can be very slow. Rewrite with a range bucket or an equality part when you can.

Hints and the planner

df.join(other.hint("merge"), "id")        # sort-merge
df.join(other.hint("shuffle_hash"), "id") # shuffle hash

The hint names are BROADCAST, MERGE, SHUFFLE_HASH and SHUFFLE_REPLICATE_NL.

With adaptive query execution, Spark can change its mind at runtime. If after a shuffle it sees that one side is actually small, it can switch a sort-merge join to a broadcast join. It can also split skewed partitions. So the planned strategy in explain() before the run can differ from what ran, and the final plan is shown in the SQL tab.

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext