What you will be able to do
- Explain how AQE coalesces small post-shuffle partitions, and why it can only reduce the partition count
- Describe when AQE switches a sort merge join to a broadcast hash join at runtime, and what that costs
- Apply the two-part rule AQE uses to decide whether a partition is skewed
- Summarize the benefits of AQE, including empty-relation propagation and reduced manual tuning
1.Dynamically coalescing shuffle partitions
Adaptive Query Execution (AQE) is Spark's runtime re-optimization of a query plan, based on statistics gathered as the query runs. The Databricks documentation lists four major AQE features, and coalescing shuffle partitions is the first problem many teams run into.
A shuffle moves data across the network so that downstream operators receive it in the layout they need. How many partitions the shuffle produces matters a lot, and the right number depends on the data, which can vary widely between stages and between queries. With too few partitions, each one is large and its task may spill to disk during a sort or aggregation. With too many, each one is tiny, and the query does many small network fetches and puts extra load on the task scheduler. The default spark.sql.shuffle.partitions value of 200 is a fixed guess.
With AQE you can start with a relatively large number of shuffle partitions and let AQE merge adjacent small ones at runtime, based on the shuffle file statistics. The Databricks documentation explains why this helps: very small tasks have worse I/O throughput and pay proportionally more for scheduling and task setup, so combining them saves resources and improves cluster throughput.
The Databricks blog that introduced AQE uses SELECT max(i) FROM tbl GROUP BY j as an example, with an initial shuffle partition number of five. Without AQE, Spark runs five final-aggregation tasks even though three of the partitions are very small. AQE merges those three, so the aggregation runs three tasks instead of five.
Coalescing only merges partitions. It never creates more partitions than the shuffle started with, which is why the Apache Spark guide tells you to set a large enough initial number (spark.sql.adaptive.coalescePartitions.initialPartitionNum, which defaults to spark.sql.shuffle.partitions). The Databricks settings that control the target size are spark.sql.adaptive.advisoryPartitionSizeInBytes (64MB; coalesced partitions will be close to but no bigger than it) and spark.sql.adaptive.coalescePartitions.minPartitionSize (1MB). Databricks also lets you set spark.sql.shuffle.partitions to auto, which turns on auto-optimized shuffle and picks the number from the query plan and the input data size.
Checkpoint 1 of 5· Check yourself
Why does the Apache Spark guide tell you to set a large enough initial number of shuffle partitions when relying on AQE coalescing?
Coalescing only combines contiguous small partitions, so the initial count is a ceiling. Starting high gives AQE room to settle on the right number at runtime.
“Spark can pick the proper shuffle partition number at runtime once you set a large enough initial number of shuffle partitions”Source: spark.apache.org
Checkpoint 2 of 5· Exam question
A PySpark job performs a wide aggregation: ```python result = (df.groupBy("region") .agg(F.sum("amount").alias("total"))) result.write.mode("overwrite").saveAsTable("agg_totals") ``` The `region` column has very few distinct values, so many of the 200 shuffle partitions created by the aggregation end up nearly empty. With AQE enabled and default settings, what happens to these shuffle partitions before the write stage runs?
Correct answer: A — AQE merges the small post-shuffle partitions into fewer, larger partitions sized toward the `spark.sql.adaptive.advisoryPartitionSizeInBytes` target before the next stage starts.
- A. This is correct: AQE's coalescing feature reads the actual size of each post-shuffle partition and merges the small, nearly empty ones together so the write stage runs with fewer, more evenly sized tasks.
- B. This is incorrect because `spark.sql.shuffle.partitions` only sets the initial partition count; AQE is specifically designed to adjust that count downward at runtime once real sizes are known.
- C. This is incorrect because AQE does not simply warn and proceed with the original layout; its coalescing logic actively merges undersized partitions rather than leaving them to run as separate near-empty tasks.
- D. This describes partial aggregation, a separate optimization for reducing shuffle volume before the exchange; it does not describe what AQE does to the partitions after the shuffle has already happened.
- E. This is incorrect because AQE coalesces toward a target partition size, not down to a single partition, and it does not force a single output file regardless of data volume.
2.Dynamically switching join strategies
Broadcast hash join is usually the fastest join when one side fits comfortably in memory. Spark plans one when the estimated size of a join side is below the broadcast threshold. That estimate can be wrong, for example after a very selective filter, or when the join side is the output of several complex operators rather than a simple scan.
AQE re-plans the join using the actual size of that side at runtime. If the measured size is below the adaptive broadcast threshold, the planned sort merge join becomes a broadcast hash join. On Databricks that threshold is spark.databricks.adaptive.autoBroadcastJoinThreshold, which defaults to 30MB. Apache Spark uses spark.sql.adaptive.autoBroadcastJoinThreshold, which defaults to the same value as spark.sql.autoBroadcastJoinThreshold.
This runtime switch is not free. Both sides have usually already been shuffled by the time AQE sees the statistics. Switching still avoids sorting both sides, and with spark.sql.adaptive.localShuffleReader.enabled Spark can read the shuffle files locally to cut network traffic. Even so, the Apache Spark guide says it is less efficient than planning a broadcast join from the start. Apache Spark can also turn a sort merge join into a shuffled hash join when every post-shuffle partition is below spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold.
Checkpoint 3 of 5· Check yourself
Halfway through a query, AQE converts a sort merge join into a broadcast hash join. Which statement about this conversion is accurate?
By the time AQE switches, the shuffle work has already been done. The switch saves the sorts and network traffic, but it can't recover what a broadcast planned from the start would have saved.
“This is not as efficient as planning a broadcast hash join in the first place”Source: spark.apache.org
3.Dynamically optimizing skew joins
Data skew means data is unevenly spread across partitions. In a join, one oversized partition becomes a straggler task that the whole stage has to wait for. Before AQE, you usually needed hints to deal with this. AQE detects skew from the shuffle file statistics, splits each skewed partition into smaller sub-partitions, and joins each sub-partition with the matching partition from the other side, replicating that partition where necessary. On Databricks this works for both sort merge join and shuffle hash join.
The blog's example joins table A to table B, where partition A0 is much larger than the rest. Without the optimization, four tasks run the join and one of them takes far longer than the others. With it, five tasks run, and they all take about the same time.
A partition counts as skewed only if it meets both conditions:
1. Its size is greater than spark.sql.adaptive.skewJoin.skewedPartitionFactor (default 5) multiplied by the median partition size.
2. Its size is greater than spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes (default 256MB).
In the prediction, 600MB is more than 5 × 100MB = 500MB and more than 256MB, so it is skewed. A 200MB partition with a 10MB median passes the factor test but fails the 256MB threshold, so AQE leaves it alone. The feature as a whole is controlled by spark.sql.adaptive.skewJoin.enabled, which defaults to true.
Checkpoint 4 of 5· Check yourself
With default settings, which partition will AQE split as skewed?
1GB is more than 5 × 150MB = 750MB and more than 256MB. The 200MB partition fails the threshold, the 300MB one fails the factor test, and the threshold alone is never enough.
“A partition is considered skewed when both (partition size > skewedPartitionFactor * median partition size) and (partition size > skewedPartitionThresholdInBytes) are true.”Source: docs.databricks.com
4.Empty relations, the settings, and what AQE buys you
The 2020 blog describes three AQE features for Spark 3.0. The current Databricks documentation lists a fourth: AQE dynamically detects and propagates empty relations. When a stage turns out to produce no rows, AQE can replace part or all of the plan with a LocalTableScan whose relation is empty, so the work that would have consumed that empty input is skipped. This feature is controlled by spark.databricks.adaptive.emptyRelationPropagation.enabled, which defaults to true.
Each feature has its own setting, and all of them are on by default:
| Property | Controls | Default |
|---|---|---|
| spark.sql.adaptive.coalescePartitions.enabled | Post-shuffle partition coalescing | true |
| spark.sql.adaptive.advisoryPartitionSizeInBytes | Target size after coalescing | 64MB |
| spark.databricks.adaptive.autoBroadcastJoinThreshold | Threshold for switching to broadcast join at runtime | 30MB |
| spark.sql.adaptive.skewJoin.enabled | Skew join handling | true |
| spark.databricks.adaptive.emptyRelationPropagation.enabled | Dynamic empty relation propagation | true |
In its TPC-DS experiments, Databricks reported speedups of up to 8x, with 32 queries running more than 1.1x faster. Most of those gains came from coalescing and join switching, because randomly generated TPC-DS data has no skew. Databricks reports bigger gains on production workloads that use all three of the original features.
The bigger benefit is in how you work. Cost-based optimization has always had to balance the cost of collecting statistics against the accuracy of its estimates. Detailed statistics such as column histograms are expensive to collect and quickly go out of date. AQE makes planning much less dependent on those statistics and on manual tuning of things like partition counts. It also copes better with arbitrary UDFs and with sudden changes in data size or skew, because it measures the data as the query runs instead of relying on what you knew beforehand.
Checkpoint 5 of 5· Exam question
A team wants to tune AQE's post-shuffle partition coalescing so that Spark targets roughly 128 MB per coalesced partition instead of the 64 MB default, without changing anything about join or skew handling. Which two configuration properties are directly involved in this behavior? (Choose 2 answers)(Select 2)
Correct answers: A, B — `spark.sql.adaptive.advisoryPartitionSizeInBytes`, which sets the target size Spark aims for when merging small post-shuffle partitions together.; `spark.sql.adaptive.coalescePartitions.enabled`, which must stay set to `true` for Spark to merge any undersized post-shuffle partitions together at all.
- A. This is correct because this property is exactly the target size AQE's coalescing logic aims for when combining small post-shuffle partitions, so raising it to 128 MB directly changes the coalescing behavior described.
- B. This is correct because coalescing only runs when this flag is `true`; it must stay enabled for the advisory partition size to have any effect on how partitions are merged.
- C. This threshold governs skew detection for splitting oversized partitions, a separate feature from coalescing, so it does not affect how small partitions get merged together.
- D. This threshold controls whether a join side is small enough to broadcast, which is unrelated to how post-shuffle partitions are sized or merged after an aggregation.
- E. This sets the initial shuffle partition count before AQE runs, but it is not the mechanism that controls the coalescing target size, and AQE can still reduce the count below it at runtime.
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.AQE coalescing can increase the number of shuffle partitions when partitions are too large.Why is that wrong?
Coalescing only combines small partitions, so the initial shuffle partition count is a ceiling. That is why you set it large enough. Only skew handling splits partitions, and it splits individual skewed ones.
Covered in Dynamically coalescing shuffle partitions
2.Because AQE switches to broadcast joins at runtime, you never need to plan a broadcast join up front.Why is that wrong?
The runtime switch happens after the shuffle has already been done. It is better than finishing the sort merge join, but less efficient than a broadcast hash join planned from the start.
Covered in Dynamically switching join strategies
3.Any partition larger than skewedPartitionThresholdInBytes is treated as skewed.Why is that wrong?
A partition must exceed both the byte threshold and skewedPartitionFactor times the median partition size. Meeting either condition alone is not enough.
Covered in Dynamically optimizing skew joins
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.https://docs.databricks.com/aws/en/optimizations/aqeOfficial docs
“Very small tasks have worse I/O throughput and tend to suffer more from scheduling overhead and task setup overhead.”
↩︎ Dynamically coalescing shuffle partitions“The coalesced partition sizes will be close to but no bigger than this target size.”
↩︎ Dynamically coalescing shuffle partitions“The threshold to trigger switching to broadcast join at runtime.”
↩︎ Dynamically switching join strategies“Dynamically handles skew in sort merge join and shuffle hash join by splitting (and replicating if needed) skewed tasks into roughly evenly sized tasks.”
↩︎ Dynamically optimizing skew joins“part of (or entire) the plan is replaced by node LocalTableScan with the relation field as empty.”
↩︎ Empty relations, the settings, and what AQE buys you“Combining small tasks saves resources and improves cluster throughput.”
↩︎ Empty relations, the settings, and what AQE buys you“A partition is considered skewed when both (partition size > skewedPartitionFactor * median partition size) and (partition size > skewedPartitionThresholdInBytes) are true.”
↩︎ Exam trap 3“Dynamically coalesces partitions (combine small partitions into reasonably sized partitions) after shuffle exchange.”
↩︎ Prediction“A partition is considered skewed when both (partition size > skewedPartitionFactor * median partition size) and (partition size > skewedPartitionThresholdInBytes) are true.”
↩︎ Checkpoint - 2.https://www.databricks.com/blog/2020/05/29/adaptive-query-execution-speeding-up-spark-sql-at-runtime.htmlSecondary source
“AQE coalesces these three small partitions into one and, as a result, the final aggregation now only needs to perform three tasks rather than five.”
↩︎ Dynamically coalescing shuffle partitions“a number of things can make this size estimation go wrong — such as the presence of a very selective filter”
↩︎ Dynamically switching join strategies“After this optimization, there will be five tasks running the join, but each task will take roughly the same amount of time”
↩︎ Dynamically optimizing skew joins“Adaptive Query Execution yielded up to an 8x speedup in query performance”
↩︎ Empty relations, the settings, and what AQE buys you“AQE has largely eliminated the need for such statistics as well as for the manual tuning effort.”
↩︎ Empty relations, the settings, and what AQE buys you“combine adjacent small partitions into bigger partitions at runtime by looking at the shuffle file statistics”
↩︎ Exam trap 1
Also cited
- https://spark.apache.org/docs/latest/sql-performance-tuning.htmlSecondary source
“This is not as efficient as planning a broadcast hash join in the first place”
↩︎ Exam trap 2“Spark can pick the proper shuffle partition number at runtime once you set a large enough initial number of shuffle partitions”
↩︎ Checkpoint“This is not as efficient as planning a broadcast hash join in the first place”
↩︎ Checkpoint