Adaptive Query Execution
Spark plans a query before it has seen your data. AQE lets it change that plan halfway through, using statistics from the stages that already ran.
On this page
Show code in
Every code block on the page follows this.
You will learn
- When and why Spark re-plans a running query
- The three changes AQE can make, with real numbers
- Which configs matter, and their defaults
- What AQE cannot fix for you
Read first
- Jobs, stages and tasks · 3 min read
- Shuffle partitions · 3 min read
- Broadcast joins · 3 min read
Comfortable with these? Read on.
The problem it solves
The optimizer picks join strategies and partition counts from estimates. Estimates are often wrong: a filter removes 99% of rows, or one key holds half the data. A plan fixed up front then shuffles into 200 nearly empty partitions, or sends one giant partition to a single task.
How it works
At every shuffle boundary, Spark pauses, measures the shuffle output, and re-optimizes the rest of the plan.
-
Step 1
A stage runs and writes its shuffle
-
Step 2
Spark reads the real partition sizes
-
Step 3
The remaining plan is re-optimized
Re-optimizing can change three things:
Coalesce partitions
Switch join strategy
Split skewed partitions
Walk through one query
Daily revenue for one country, from a 500 GB orders table joined to customers:
SELECT o.order_date, SUM(o.amount) FROM orders o JOIN customers c ON o.customer_id = c.id WHERE c.country = 'NZ' GROUP BY o.order_date
| Decision | Planned up front | What the stats showed | After AQE |
|---|---|---|---|
| Join | Sort-merge: the filtered customers size was guessed as large | Filtered customers are 6 MB | Broadcast hash join (under the 10 MB threshold) |
| Aggregation partitions | 200, from spark.sql.shuffle.partitions | Only 3 GB shuffled: 200 partitions averaging 15 MB, many far smaller | About 47 partitions of roughly 64 MB |
The arithmetic behind the second row: 3 GB divided by the 64 MB advisory size is about 47. Fewer, fuller tasks mean less scheduling overhead and fewer tiny output files, with no change to your code. When AQE converts a join to broadcast after the shuffle already ran, it can also read the shuffle files locally on each executor (a "local shuffle reader") instead of fetching them over the network.
Key configs
| Config | Default | What it does |
|---|---|---|
spark.sql.adaptive.enabled | true | Turns AQE on or off. |
…coalescePartitions.enabled | true | Merges small shuffle partitions. |
…advisoryPartitionSizeInBytes | 64MB | Target partition size when merging or splitting. |
…autoBroadcastJoinThreshold | same as the non-adaptive one | Runtime size below which a join side is broadcast. |
…skewJoin.enabled | true | Splits skewed partitions in sort-merge joins. |
spark.conf.set("spark.sql.adaptive.enabled", "true") spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128MB")
SET spark.sql.adaptive.enabled = true; SET spark.sql.adaptive.advisoryPartitionSizeInBytes = 128MB;
Spot it in a plan
A plan that AQE manages starts with AdaptiveSparkPlan. Before an action runs it shows isFinalPlan=false; after it runs, the SQL tab of the Spark UI shows the final, re-optimized plan, with nodes such as AQEShuffleRead where partitions were coalesced or split.
== Physical Plan == AdaptiveSparkPlan isFinalPlan=false +- SortMergeJoin [customer_id], [id], Inner ...
When AQE cannot help
AQE only acts at shuffle boundaries, using statistics it has already measured. That leaves real gaps:
- The first stage. Nothing has been measured yet, so the file scan and anything before the first shuffle run as planned.
- Queries with no shuffle. A filter-and-write job has no boundary to re-plan at.
- Skewed aggregations. Skew handling splits partitions in sort-merge joins. A hot key in a
groupBystill lands on one task. See Data skew and salting. - Structured Streaming. Streaming queries do not use AQE.
Common mistakes
Setting shuffle partitions low and expecting AQE to raise them
Reading explain() and treating it as the final plan
isFinalPlan=false. The join you see may become a broadcast at runtime. Check the SQL tab in the Spark UI after the job runs.Turning AQE off to debug, and leaving it off
Key takeaways
- AQE re-plans at each shuffle boundary using measured sizes, not estimates.
- It coalesces small partitions, switches joins to broadcast, and splits skewed join partitions.
- It is on by default since Spark 3.2. Set shuffle partitions high and let it merge down.
- It cannot help the first stage, shuffle-free queries, skewed aggregations or streaming.
Check yourself
3 questions1. A query shuffles into 200 partitions, but most hold only a few kilobytes. AQE is on. What happens?
Show the answer
Small partitions are merged toward the advisory size, so fewer tasks run. AQE coalesces post-shuffle partitions toward spark.sql.adaptive.advisoryPartitionSizeInBytes (64 MB by default).
2. Which of these can AQE NOT fix?
Show the answer
A hot key in a groupBy aggregation. Skew handling applies to sort-merge joins, not aggregations.
3. What does isFinalPlan=false in explain() mean?
Show the answer
The plan may still change at runtime. AQE re-optimises during execution, so the plan before the action is provisional.
Practice it
Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.
Go deeper
AQE skew-join handlingSpark internals
Shuffle partitionsSpark internals
Broadcast joins
Primary sources: Adaptive Query Execution