CertSafari
    Databricks Certified Associate Developer for Apache Spark· Lessons

    Domain 4 · Lesson 23/32

    Adaptive Query Execution: How Spark Re-optimizes Queries at Runtime

    Describe Adaptive Query Execution (AQE) and its benefits.

    10 min read
    3.12% of exam
    3 sources
    Published 3 Oct 2026
    Docs as of 30 Sep 2026

    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?

    Sources12

    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. 1.Re-run the optimizer, physical planner and adaptive physical rules on the updated plan
    2. 2.Kick off all leaf stages, the stages that depend on no other stage
    3. 3.Repeat the cycle until the entire query is done
    4. 4.As stages finish materializing, update the logical plan with their runtime statistics
    5. 5.Execute new query stages whose child stages have all materialized

    Sources12

    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.

    Umbrella switches for AQE as each documentation set lists them
    PropertyDocumented inDefault
    spark.databricks.optimizer.adaptive.enabledDatabricks AQE docstrue
    spark.sql.adaptive.enabledApache Spark performance tuning guidetrue (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?

    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?

    Sources132

    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.

    Plan signatures that show an AQE optimization took effect
    AQE optimizationWhat to look for in the current/final plan
    Sort merge join changed to broadcast hash joinA different physical join node from the one in the initial plan
    Coalesced partitionsNode CustomShuffleReader with property Coalesced
    Skew join handlingNode SortMergeJoin with field isSkew as true
    Empty relation propagationPart 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.

    Sources1

    Exam traps

    Each one states something that sounds right. Open it to see what is actually true.

    1. 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. 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. 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. 1.
      “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. 2.
      “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. 3.

    Continue to page 2 of 2

    AQE Features: Partition Coalescing, Join Switching and Skew Join Handling

    Spotted a mistake, or was something unclear? Tell us.