In Spark AQE, how can a sort-merge join become a broadcast hash join mid-query?
answer
- the planner guesses; the shuffle knows
- stages are re-optimised as they complete
- a filtered side turns out tiny after all
- the shuffle is already written by then
- check which plan was final, not initial
basics
~20 sAdaptive Query Execution re-plans at shuffle boundaries. Once a join side's shuffle has been written, Spark knows its exact size; if that is under the adaptive broadcast threshold, it replaces the planned sort-merge join with a broadcast hash join.
solid answer
~40 sAQE splits the physical plan into query stages at each shuffle. When a stage finishes, its map-output statistics give exact sizes rather than the planner's estimates — which is where the original decision usually went wrong, because a side that looked large before a selective filter or an upstream join turns out to be tiny. If a side is smaller than `spark.sql.adaptive.autoBroadcastJoinThreshold` (which defaults to the value of `spark.sql.autoBroadcastJoinThreshold`), Spark re-plans the join as a broadcast hash join. This is worse than planning a broadcast join up front, since both shuffles were already written, but it avoids sorting and merging both sides. It also enables the local shuffle reader (`spark.sql.adaptive.localShuffleReader.enabled`, true by default), so executors read shuffle files locally instead of fetching them over the network.
code
properties · 4 linesspark.sql.adaptive.enabled=true
# defaults to the value of spark.sql.autoBroadcastJoinThreshold
spark.sql.adaptive.autoBroadcastJoinThreshold=10m
spark.sql.adaptive.localShuffleReader.enabled=truego deeper
Know that Spark 3 can change the join strategy while a query is running, because after a shuffle it knows the real data sizes instead of the estimates it started from.
Explain query stages and where the runtime statistics come from, and name the adaptive broadcast threshold and its relationship to the plan-time one.
Use it deliberately: recognise when a runtime conversion is masking missing table statistics, and prefer fixing the estimate or hinting the broadcast so the large side's shuffle is never written.
Decide how much the platform should lean on adaptive re-planning versus maintained statistics, and what that means for cost predictability when a plan can change between two runs of the same nightly job.
## Query stages and re-planning Without adaptive execution, Spark chooses every physical operator before the first task runs, using estimated row counts and sizes derived from data-source metadata, catalog statistics and the optimizer's selectivity guesses. Those estimates decay fast: after a filter on an un-analyzed column, after a join whose output cardinality nobody knows, after a UDF, the number the planner is holding may be off by orders of magnitude. AQE turns the plan into a sequence of *query stages* separated by shuffles. It executes a stage, reads its real map-output statistics, re-optimises the remainder of the plan with those exact numbers, and continues. The join-strategy switch is the most visible consequence. ## The conversion When the runtime statistics of any side of a sort-merge join come out below the adaptive broadcast threshold, AQE rewrites that join into a broadcast hash join for the stages that have not yet run. `spark.sql.adaptive.autoBroadcastJoinThreshold` controls it and defaults to the same value as `spark.sql.autoBroadcastJoinThreshold` (10 MB); setting it to `-1` disables adaptive broadcasting specifically, without touching the plan-time threshold. The documentation is explicit that this is not as good as planning a broadcast hash join in the first place — both sides have already paid the shuffle write — but better than continuing with the sort-merge join, because the sort of both sides and the merge are avoided. ## The local shuffle reader Once the join no longer needs data partitioned by key, the reducer tasks do not have to fetch a specific bucket from every executor. With `spark.sql.adaptive.localShuffleReader.enabled` true (the default), Spark reads the shuffle output that already sits on the local executor, cutting the network fetch out of the picture. This is a small optimisation that only makes sense in exactly this situation: the shuffle exists, but its partitioning has become irrelevant. ## A second conversion AQE can also convert a sort-merge join into a **shuffled hash join** — building a hash map per partition instead of sorting both sides — when all post-shuffle partitions are below `spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold`. That config defaults to `0`, so this conversion is effectively off until you set a real per-partition size. It is worth knowing about mostly as evidence that AQE's join re-planning is not limited to broadcasting. ## Why this matters in practice The conversion is the reason a Spark 3 job can be dramatically faster than the same code on Spark 2 with no changes: the pathological case of a sort-merge join between a huge table and a tiny filtered one, which the estimator failed to spot, now fixes itself after the shuffle. It is also the reason "my plan changed halfway through the query" is a normal observation rather than a bug report. The practical consequence for tuning: if a broadcast is what you want, getting it at *plan* time is still better. Give the planner real statistics with `ANALYZE TABLE`, or state your intent with a `broadcast()` hint, and you skip the shuffle entirely instead of paying for it and then discarding its partitioning. ## Seeing it happen An adaptive plan is rooted at an `AdaptiveSparkPlan` node carrying an `isFinalPlan` flag. Before execution completes it reads `isFinalPlan=false` and shows the initial plan; the SQL tab in the Spark UI updates the plan as stages complete, and the final plan — the one that actually ran — is what you should compare against your expectations. Reading `explain()` output on an AQE query and concluding that a sort-merge join ran, when in fact the final plan broadcast, is a routine mistake.
- If AQE can broadcast at runtime, why still collect statistics or use a broadcast hint?Because the runtime switch happens after both sides have been shuffled. A plan-time broadcast skips the exchange on the large side entirely, which is usually the dominant cost. Running ANALYZE TABLE or adding an explicit broadcast hint gets you the cheaper plan; the adaptive conversion is the safety net for cases the estimator could not have known.
- How do you tell which plan actually executed?Look at the final plan, not the initial one. An adaptive query's root node reports isFinalPlan, and the Spark UI's SQL tab rewrites the plan as query stages complete. Reading explain() output before execution shows the plan AQE started with, which may not be the join strategy that ran.
- Can AQE also switch a sort-merge join to a shuffled hash join?Yes, when every post-shuffle partition is small enough to build a local hash map, governed by spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold. That threshold defaults to 0, so the conversion is effectively disabled until you set a real per-partition size — worth knowing about, but not something you meet with stock settings.
saying these in an interview costs you the question
- Thinks AQE re-reads the source to get better statistics
- Believes the runtime broadcast is as cheap as a planned one
- Says the adaptive threshold is unrelated to the plan-time one
- Reads the initial explain output as the plan that executed
- Assumes a plan cannot change once the query has started