Choosing Between Broadcast and Sort-Merge Joins in Apache Spark
A decision guide for picking the right join strategy on Spark DataFrames, with a comparison table, trade‑off analysis, and a PySpark example you can validate in the Spark UI.
17 Oct 2025, 09:05 UTC

Decision and Constraints
When joining two Spark DataFrames on an equality condition you must decide whether to let Catalyst choose the physical join (usually a sort‑merge join for large tables) or to force a broadcast hash join when one side is small enough to fit in executor memory. The decision hinges on four constraints:
- Table size – the estimated size of the build side (the side that will be broadcast or hashed).
- Executor memory – available memory per executor after accounting for overhead, caching, and other operations.
- Data skew – whether keys are unevenly distributed, which can cause stragglers in shuffle‑based joins.
- Shuffle tolerance – how much network I/O and disk spill you are willing to accept for the join.
If the build side is small (< broadcast threshold) and skew is low, a broadcast hash join eliminates shuffle on the large side and is usually fastest. If both sides are large or skew is high, a sort‑merge join is more robust. A shuffle hash join can be useful when each partition of the build side fits in memory but you want to avoid the sort cost of sort‑merge.
Options Comparison
| Option | When Optimizer Picks It | Shuffle Volume | Memory Risk | Skew Behavior | How to Force |
|---|---|---|---|---|---|
| Broadcast Hash Join | One side estimated < spark.sql.autoBroadcastJoinThreshold (default ~10 MiB) and AQE does not demote it. |
None on the large side; only the small side is sent to each executor (driver‑to‑executor broadcast). | High if the small side exceeds executor memory or driver memory (OOM or broadcast timeout). | Resilient to skew because the large side is not shuffled; only the small side is replicated. | Wrap the small side with broadcast(df) in PySpark or use the /*+ BROADCAST(df) */ hint in SQL; raise spark.sql.autoBroadcastJoinThreshold to increase the limit. |
| Sort‑Merge Join | Both sides are large (> broadcast threshold) or statistics are missing; default for equi‑joins in Spark 3.x. | Full shuffle of both sides (shuffle read + write) plus external sort if data exceeds memory. | Low to moderate; each partition only needs enough memory to hold a slice during the merge step. | AQE skew‑join handling (Spark 3.2+) can split skewed partitions to reduce stragglers. | Disable broadcast globally (spark.sql.autoBroadcastJoinThreshold=0) or use /*+ MERGE */ hint. |
| Shuffle Hash Join | Enabled when spark.sql.join.preferSortMergeJoin=false and each build‑side partition fits in memory. |
Full shuffle of both sides (no sort). | Medium; each executor must hold an in‑memory hash map for its partition of the build side. | Similar to sort‑merge; skewed partitions can still cause hot executors. | Set spark.sql.join.preferSortMergeJoin=false and optionally increase spark.sql.shuffle.partitions to keep partition size small. |
Trade‑off Summary
- Broadcast wins when you have a star‑schema fact table joining a small dimension (e.g., a lookup table of < 100 MiB). Shuffle on the fact side is avoided, giving the lowest latency.
- Sort‑merge wins for large‑to‑large joins or when you cannot guarantee the small side fits in memory; it spills to disk and works at any scale.
- Shuffle hash is a niche middle ground: useful when you want to avoid the sort cost but still need a shuffle, and you can tune partition size so each build‑side partition stays in memory.
Accurate table statistics are essential. For file‑based sources (Parquet, ORC, etc.) run ANALYZE TABLE … COMPUTE STATISTICS or cache the small side so its size is known; otherwise the optimizer may miss a broadcast opportunity.
Concrete Implementation (PySpark)
The following snippet demonstrates forcing a broadcast join, verifying the plan, and checking shuffle metrics in the Spark UI.
from pyspark.sql import SparkSession
from pyspark.sql.functions import broadcast
spark = SparkSession.builder \
.appName("JoinStrategyDemo") \
.getOrCreate()
# Assume fact_df is large (e.g., sales events) and dim_df is small (e.g., product catalog)
fact_df = spark.read.parquet("s3://my-bucket/fact/")
dim_df = spark.read.parquet("s3://my-bucket/dim/")
# Force broadcast of the small side
joined = fact_df.join(broadcast(dim_df), "product_key")
# Show the logical and physical plan before execution
print("=== Explain (formatted) ===")
joined.explain(True)
# Trigger an action to materialize the join
joined.count() # or .write.mode("overwrite").parquet("s3://my-bucket/output/")
After the .count() action completes, open the Spark UI:
- Navigate to the SQL tab for the application.
- Find the most recent query (look for the join operation).
- In the Details pane, verify that the Physical Plan contains
BroadcastHashJoin(notSortMergeJoin). - Check the Shuffle Read and Shuffle Write metrics; they should be near zero for the fact side, confirming that only the small side was broadcast.
To compare with the default sort‑merge join, repeat the same code without broadcast() (or set spark.sql.autoBroadcastJoinThreshold=0) and observe the plan change to SortMergeJoin with measurable shuffle read/write bytes.
Limitations and Practical Checks
- Broadcast size estimate – The optimizer uses statistics; if they are stale or missing, it may underestimate the size. Always run
ANALYZE TABLEor cache the small side before relying on a broadcast. - Driver memory pressure – The small side is collected to the driver before broadcasting. Very wide rows or many concurrent broadcasts can exhaust driver memory even if each table fits in executor memory. Monitor driver GC and memory usage in the UI.
- AQE runtime changes – With Adaptive Query Execution enabled (default in Spark 3.x), the plan shown by
explain()can change at runtime (e.g., a broadcast may be demoted if the shuffled size exceeds expectations). Validate the final plan in the SQL tab after execution, not just the pre‑explain output. - Non‑equi joins – Broadcast hints only apply to equi‑joins. For complex conditions Spark may fall back to
BroadcastNestedLoopJoinor a cartesian product, which do not scale. Verify the plan showsBroadcastHashJoinfor your equality condition.
How to check the result – After running the join:
- Run
df.explain('formatted')and confirm the stringBroadcastHashJoinappears where expected. - In the Spark UI SQL tab, locate the query and confirm Shuffle Read Bytes for the large side is close to zero (e.g., < 1 MiB).
- Optionally, capture the execution time with a simple wrapper (
time.time()before and after the.count()) and compare it to the sort‑merge baseline.
Summary
Choose a broadcast hash join when one side is small enough to fit in executor memory, statistics are reliable, and you can tolerate the driver‑side memory cost. Otherwise let Spark pick the default sort‑merge join, which handles any scale and skew gracefully. Use the table above to weigh shuffle volume, memory risk, and skew behavior, then validate the decision with explain() and the Spark UI.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.