Skip to content
Good engineers know 4 min read · spark.sql.adaptive.enabled 2 practice problems ↓

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.

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

Comfortable with these? Read on.

TL;DR At every shuffle boundary, AQE measures the real shuffle output and re-optimises the rest of the query: it merges small partitions, switches joins to broadcast, and splits skewed join partitions. On by default since Spark 3.2.

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.

  1. Step 1

    A stage runs and writes its shuffle

  2. Step 2

    Spark reads the real partition sizes

  3. Step 3

    The remaining plan is re-optimized

Re-optimizing can change three things:

Coalesce partitions

Merges many tiny shuffle partitions into fewer, right-sized ones.

Switch join strategy

Turns a sort-merge join into a broadcast join when one side turns out small.

Split skewed partitions

Breaks an oversized join partition into several tasks.

Walk through one query

Daily revenue for one country, from a 500 GB orders table joined to customers:

Spark SQL
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
DecisionPlanned up frontWhat the stats showedAfter AQE
JoinSort-merge: the filtered customers size was guessed as largeFiltered customers are 6 MBBroadcast hash join (under the 10 MB threshold)
Aggregation partitions200, from spark.sql.shuffle.partitionsOnly 3 GB shuffled: 200 partitions averaging 15 MB, many far smallerAbout 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

ConfigDefaultWhat it does
spark.sql.adaptive.enabledtrueTurns AQE on or off.
…coalescePartitions.enabledtrueMerges small shuffle partitions.
…advisoryPartitionSizeInBytes64MBTarget partition size when merging or splitting.
…autoBroadcastJoinThresholdsame as the non-adaptive oneRuntime size below which a join side is broadcast.
…skewJoin.enabledtrueSplits skewed partitions in sort-merge joins.
PySparkSpark SQL · Setting configs
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 groupBy still 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

Coalescing only merges. If you start with 20 partitions for 2 TB, each one is about 100 GB and AQE will not split them. Start high and let AQE merge down.

Reading explain() and treating it as the final plan

Before an action, the plan says 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

A common way for jobs to get mysteriously slower after an incident. Turn it off for one session only, never in cluster defaults.

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 questions

1. 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.

Solve: Pick the Join Strategy →

Go deeper

Primary sources: Adaptive Query Execution