Spark AQE in Practice: When the Runtime Should Pick Your Shuffle Partitions
Static planning guesses at shuffle partition counts and join strategies. AQE re-decides after the shuffle, using real partition sizes — here's how to configure and verify it.
23 Jan 2026, 03:19 UTC

Spark's static planner chooses a shuffle partition count before any data moves. spark.sql.shuffle.partitions defaults to 200, so a job shuffling 2 GB and a job shuffling 2 TB both start with 200 partitions unless someone overrides it. The failure mode is familiar: either hundreds of tiny tasks that spend more time on scheduling than on work, or a handful of oversized tasks that dominate the stage while the rest of the cluster idles.
Adaptive Query Execution (AQE) changes when those decisions are made. Rather than committing to a partition count and join strategy at planning time, Spark re-optimizes the plan at runtime, after shuffle map stages have produced real statistics. This post covers what AQE actually re-decides, a worked configuration, how to verify it took effect, and where it stops helping.
What AQE re-decides, and when
AQE operates in units called query stages, which are separated by shuffle boundaries. Once a shuffle map stage finishes, Spark knows the actual size of each partition rather than an estimate derived from table statistics. It then applies three rewrites:
- Coalescing shuffle partitions. Many small partitions are merged into fewer, right-sized ones. The target is controlled by
spark.sql.adaptive.advisoryPartitionSizeInBytes. - Switching join strategies. A sort-merge join planned before the shuffle can be replaced by a broadcast join if one side turns out to be small enough at runtime.
- Splitting skewed partitions. For sort-merge joins, a partition much larger than its peers can be split into sub-partitions, with the matching side replicated to each.
Two constraints matter. AQE only re-optimizes across shuffle boundaries, so a query with no shuffle gives it nothing to adapt to. And AQE lives in the Catalyst optimizer, so it applies to Spark SQL, DataFrames, and Datasets — not to the low-level RDD API.
A worked example: group-by after a shuffle
Consider a DataFrame aggregation that groups by a moderately selective column and then joins the result to a small dimension table. The shuffle is unavoidable, the output partitions are uneven, and the dimension table is small enough to broadcast but the planner's estimate says otherwise.
Run the following from the driver, either as spark-submit arguments or in spark-defaults.conf. Per-application --conf flags need no special permissions; editing cluster-wide spark-defaults.conf requires administrator access on the cluster.
spark-submit \
--conf spark.sql.adaptive.enabled=true \
--conf spark.sql.adaptive.coalescePartitions.enabled=true \
--conf spark.sql.adaptive.advisoryPartitionSizeInBytes=64m \
--conf spark.sql.adaptive.skewJoin.enabled=true \
--conf spark.sql.shuffle.partitions=400 \
app.pyProperty names and defaults have shifted across Spark 3.x releases. AQE is enabled by default from Spark 3.2 onward; Spark 3.0 and 3.1 require the explicit opt-in shown above. Check your version's configuration documentation before assuming a given sub-feature is on.
Note the deliberately high spark.sql.shuffle.partitions value. AQE coalesces partitions downward but does not create more than the initial count, so setting this too low removes the room AQE needs to work with. Leaving it at a generous value and letting AQE merge is usually safer than hand-tuning it per job.
Verifying that AQE actually ran
- Open the Spark UI and go to the SQL tab. Select the query and inspect the physical plan. An
AdaptiveSparkPlannode indicates AQE is in play; the final plan shows the partition count chosen after runtime statistics arrived. - Check the Stages tab for shuffle read and write sizes and the distribution of task durations. Coalescing shows up as fewer, larger tasks; skew splitting shows up as a previously dominant task breaking into several.
- Run
df.explain("formatted")from the driver to see the plan structure. Be aware thatEXPLAINshows the plan before runtime statistics are applied, so it confirms AQE is configured but not what it decided. - Compare end-to-end duration and resource usage with AQE on and off on a representative dataset. The effect is workload-dependent; a single small run proves little.
Where AQE stops helping
AQE is a runtime correction, not a substitute for data layout. It cannot add partition pruning that the query never expressed, cannot merge a directory of tiny files, and cannot fix a schema that forces a shuffle where none was needed. If the physical layout is wrong, AQE will re-optimize a plan that was already doomed.
Skew handling has a hard limit too: a single hot key that cannot be divided further stays in one sub-partition, so the largest indivisible key remains the bottleneck. AQE also will not always choose a broadcast join — that depends on runtime size estimates and the configured broadcast threshold, which you may need to adjust deliberately.
Finally, runtime re-optimization costs a planning round-trip per query stage. For very small or latency-sensitive queries, that overhead can outweigh the benefit. If a job runs in well under a second, measure before enabling AQE by default.
What to do next
Enable AQE, leave spark.sql.shuffle.partitions at a generous value, and read the SQL tab rather than guessing. Tune advisoryPartitionSizeInBytes and the skew thresholds against your own workload, and treat any partition count copied from a blog post — including this one — as a starting point rather than a setting. The check that matters is the one on your data: same query, AQE on and off, compared in the Spark UI and on wall-clock time.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.