What you will be able to do
- Explain what Adaptive Query Execution (AQE) is and why it re-optimizes a query while the query runs instead of only before it starts
- Describe how the execute-reoptimize-execute loop uses query stages and runtime statistics
- State when AQE applies to a query and how it is turned on or off
- Read AdaptiveSparkPlan, isFinalPlan and isRuntime in a query plan to tell whether AQE changed it
Key concept
Adaptive Query Execution (AQE) — AQE re-optimizes a query's physical plan while the query is running. At each query-stage boundary (a shuffle or broadcast exchange) it replaces compile-time estimates with real runtime statistics and uses them to plan the rest of the query.
1.Why plan a query twice?
Spark's optimizer chooses a physical plan before any data moves. To do that it relies on statistics such as row counts, distinct values and min/max values, and the cost-based optimizer uses them to pick join types, build sides and join order. The problem is that the statistics can be wrong. Outdated statistics and imperfect cardinality estimates can produce a suboptimal plan, and once that plan is running nothing can correct it.
Adaptive Query Execution (AQE) changes this. It keeps the static plan as a starting point and then revises it while the query runs, using what it has learned about the data so far.
Spark operators are normally pipelined and run in parallel, but a shuffle or broadcast exchange breaks that pipeline. Each part of the query bounded by these materialization points is called a query stage. A query stage has to write out its whole intermediate result before the next stage can start. At that moment, statistics for every partition are available and nothing downstream has run yet, which makes it the right time to re-optimize.
With these runtime statistics, Databricks can choose a better physical strategy, pick a better post-shuffle partition size and count, or handle skewed joins without a hint. The Databricks documentation names three situations where this helps most: statistics collection is turned off, the statistics are stale, or statically derived estimates are inaccurate. Estimates tend to be inaccurate in the middle of a complicated query and after data skew has occurred.
Checkpoint 1 of 5· Check yourself
Why is the boundary between two query stages a good moment for AQE to re-optimize?
A query stage materializes its full result before the next stage starts, so complete statistics exist while downstream work can still be re-planned.
“for it is when data statistics on all partitions are available and successive operations have not started yet”Source: www.databricks.com
2.The execute-reoptimize-execute loop
The 2020 Databricks engineering blog that introduced AQE describes how the framework works through a query. It starts by running the leaf stages, which are the stages that depend on no other stage. When one or more of them finish materializing, the framework marks them complete in the physical plan and updates the logical plan with the runtime statistics from those stages.
Next it runs the optimizer again with a selected list of logical rules, then the physical planner and the physical optimization rules. Those physical rules include the regular ones plus adaptive-only rules such as partition coalescing and skew join handling. The result is a newly optimized plan in which some stages are already complete. The framework then looks for stages whose child stages have all materialized, runs them, and repeats the cycle until the query finishes.
Because re-optimization only happens between stages, any part of the plan that has already run is fixed. Only the part that hasn't run yet can still change.
Checkpoint 2 of 5· Put it in order
Put the steps of the AQE framework's runtime loop in order.
- 1.Re-run the optimizer, physical planner and adaptive physical rules on the updated plan
- 2.Kick off all leaf stages, the stages that depend on no other stage
- 3.Repeat the cycle until the entire query is done
- 4.As stages finish materializing, update the logical plan with their runtime statistics
- 5.Execute new query stages whose child stages have all materialized
AQE starts with the leaf stages, collects their runtime statistics, re-plans the remaining work, runs the stages that are now ready, and loops until the query completes.
“repeat the above execute-reoptimize-execute process until the entire query is done.”Source: www.databricks.com
3.When AQE applies, and how to switch it on
The current Databricks documentation states that AQE is enabled by default. The Apache Spark tuning guide says the same for open-source Spark from version 3.2.0. In Spark 3.0, where AQE first appeared, it was off by default and had to be enabled with spark.sql.adaptive.enabled. Open-source Spark uses spark.sql.adaptive.enabled as the umbrella switch for AQE. The Databricks documentation lists its own property for the same purpose.
| Property | Documented in | Default |
|---|---|---|
| spark.databricks.optimizer.adaptive.enabled | Databricks AQE docs | true |
| spark.sql.adaptive.enabled | Apache Spark performance tuning guide | true (enabled by default since Spark 3.2.0) |
Turning AQE on does not mean every query uses it. AQE applies only to queries that meet both of these conditions:
- The query is not streaming. - The query contains at least one exchange (usually because of a join, aggregate or window), at least one sub-query, or both.
This follows from how AQE works. Without an exchange or a sub-query there is no stage boundary at which to re-optimize. And a query that AQE applies to may still run exactly as originally planned: re-optimization might produce a different plan from the static one, or it might not.
Checkpoint 3 of 5· Check yourself
AQE is enabled on a cluster. Which of these queries is AQE NOT applied to?
AQE applies only to non-streaming queries that contain an exchange or a sub-query. The streaming query fails the first condition.
“Contain at least one exchange (usually when there's a join, aggregate, or window), one sub-query, or both.”Source: docs.databricks.com
Checkpoint 4 of 5· Exam question
A data engineer wants to understand how Adaptive Query Execution (AQE) differs from Spark's traditional cost-based optimizer. Which statement correctly describes when AQE takes effect?
Correct answer: A — AQE re-optimizes the physical plan at runtime using actual statistics collected at shuffle and broadcast exchange boundaries, replacing the original pre-execution estimates.
- A. This is correct because AQE's defining behavior is re-planning at runtime: it pauses at each exchange boundary, gathers real partition and row-count statistics, and uses them to pick a better plan for the next stage than the static optimizer could estimate up front.
- B. This is incorrect because rewriting the logical plan from the parsed AST based on indexes happens, if at all, during static analysis before execution, not through AQE's runtime re-optimization mechanism.
- C. This describes Delta file compaction (such as OPTIMIZE), a storage-layer maintenance operation that is unrelated to AQE's per-query runtime re-planning of shuffle and join stages.
- D. This describes automatic caching of repeated lineage, which is not a behavior AQE performs; AQE only reacts to statistics gathered at exchange boundaries, not repeated transformation references.
- E. This describes cluster warm-up based on historical timing data, which Spark does not do; AQE's decisions come from statistics measured during the current query's own execution, not past runs.
4.Seeing AQE in a query plan
You can check whether AQE ran, and what it changed, in three places: the Spark UI, DataFrame.explain() and SQL EXPLAIN.
In a query where AQE applies, there are one or more AdaptiveSparkPlan nodes, usually at the root of each main query or sub-query. Each AdaptiveSparkPlan node has an isFinalPlan flag. It is false before the query runs and while it runs, and it becomes true once execution completes. Under the node you see the initial plan, which is the plan before any AQE optimization, and the current plan (or the final plan, once execution has completed). In the Spark UI, the plan diagram updates as the query runs. Nodes that have already executed do not change, but nodes that haven't executed yet can.
Each shuffle and broadcast stage also reports statistics. Before or during the stage these are compile-time estimates with isRuntime=false. After the stage completes they are the statistics collected at runtime, for example Statistics(sizeInBytes=658.1 KiB, rowCount=2.81E+4, isRuntime=true).
SQL EXPLAIN does not execute the query, so its current plan is always the same as the initial plan. It cannot show what AQE will eventually do.
| AQE optimization | What to look for in the current/final plan |
|---|---|
| Sort merge join changed to broadcast hash join | A different physical join node from the one in the initial plan |
| Coalesced partitions | Node CustomShuffleReader with property Coalesced |
| Skew join handling | Node SortMergeJoin with field isSkew as true |
| Empty relation propagation | Part or all of the plan replaced by LocalTableScan with an empty relation field |
Checkpoint 5 of 5· Match them up
Match each plan detail to what it tells you.
Tap a term, then the definition that fits it.
isFinalPlan and isRuntime show whether the plan or the statistics are final. The two named nodes are how coalescing and skew handling appear in the current/final plan.
“Dynamically coalesce partitions: node CustomShuffleReader with property Coalesced”Source: docs.databricks.com
SQL EXPLAIN never runs the query, so it only shows the initial plan and cannot reveal runtime changes. To see what AQE did, run the query and look at the final plan in the Spark UI or in DataFrame.explain() output after execution, where isFinalPlan is true.
Sources1
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.AQE is off unless you enable it with spark.sql.adaptive.enabled.Why is that wrong?
That was true only in Spark 3.0. Current Databricks documentation says AQE is enabled by default, and Apache Spark has enabled it by default since 3.2.0.
Covered in When AQE applies, and how to switch it on
2.If AQE applies to a query, its execution plan will be different from the statically compiled plan.Why is that wrong?
AQE applies to every non-streaming query with an exchange or sub-query, but re-optimization may end up with the same plan as the static one.
Covered in When AQE applies, and how to switch it on
3.SQL EXPLAIN shows the plan AQE will actually execute.Why is that wrong?
EXPLAIN does not run the query, so its current plan is always the initial plan. AQE's changes are visible only after execution, in the Spark UI or in DataFrame.explain().
Covered in Seeing AQE in a query plan
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
“This can be very useful when statistics collection is not turned on or when statistics are stale.”
↩︎ Why plan a query twice?“Nodes that have already been executed (in which metrics are available) will not change”
↩︎ The execute-reoptimize-execute loop“AQE is enabled by default.”
↩︎ When AQE applies, and how to switch it on“Before the query runs or when it is running, the isFinalPlan flag of the corresponding AdaptiveSparkPlan node shows as false”
↩︎ Seeing AQE in a query plan“Statistics(sizeInBytes=658.1 KiB, rowCount=2.81E+4, isRuntime=true)”
↩︎ Seeing AQE in a query plan“Adaptive query execution (AQE) is query re-optimization that occurs during query execution.”
↩︎ Key concept“AQE is enabled by default.”
↩︎ Exam trap 1“Not all AQE-applied queries are necessarily re-optimized.”
↩︎ Exam trap 2“As SQL EXPLAIN does not execute the query, the current plan is always the same as the initial plan”
↩︎ Exam trap 3“Databricks has the most up-to-date accurate statistics at the end of a shuffle and broadcast exchange (referred to as a query stage in AQE).”
↩︎ Prediction“Contain at least one exchange (usually when there's a join, aggregate, or window), one sub-query, or both.”
↩︎ Checkpoint“Dynamically coalesce partitions: node CustomShuffleReader with property Coalesced”
↩︎ Checkpoint - 2.https://www.databricks.com/blog/2020/05/29/adaptive-query-execution-speeding-up-spark-sql-at-runtime.htmlSecondary source
“a shuffle or broadcast exchange breaks this pipeline.”
↩︎ Why plan a query twice?“the Adaptive Query Execution framework first kicks off all the leaf stages — the stages that do not depend on any other stages.”
↩︎ The execute-reoptimize-execute loop“(default false in Spark 3.0)”
↩︎ When AQE applies, and how to switch it on“for it is when data statistics on all partitions are available and successive operations have not started yet”
↩︎ Checkpoint“repeat the above execute-reoptimize-execute process until the entire query is done.”
↩︎ Checkpoint - 3.https://spark.apache.org/docs/latest/sql-performance-tuning.htmlSecondary source
“which is enabled by default since Apache Spark 3.2.0”
↩︎ When AQE applies, and how to switch it on