Adaptive query execution: Spark re-plans a query while it runs, for instance turning a sort-merge join into a broadcast join, merging shuffle partitions, dealing with skew and propagating empty relations. Databricks advises leaving it enabled.
Also called adaptive query execution.
Read more: Microsoft Learn
In the Ultra Transcenders books
Each book explains AQE in context, with comparison tables and the common traps.
Terms in this definition
- Agents (classic) API
First-generation Foundry Agent Service API, based on threads, messages and runs. It is deprecated, replaced by conversations and responses, and retires on 31 March 2027.
- JOIN
Combines rows from two or more tables in one SELECT, most often by pairing a primary key with the foreign key that refers to it.
- Broadcast join
A join method that copies the smaller table to every executor so that the larger table needs no shuffle. You can force it with the
BROADCASThint (orBROADCASTJOIN/MAPJOIN), or AQE may pick it while the query runs. - Shuffle
When Spark has to move rows between executors, as joins, sorts and aggregations require. Per stage, the Spark UI lists Shuffle Read and Shuffle Write; heavy shuffling frequently causes spill.
- Skew
Unbalanced partitions leave a few tasks running long after the others finish. In Summary Metrics on the Spark UI, Max duration 50% or more higher than the 75th percentile points to it.
- Databricks
Analytics platform built on Apache Spark where data is processed in notebooks; Unity Catalog is the governance model it recommends.
Related terms
- Partitioning hints
The SQL hints
COALESCE,REPARTITION,REPARTITION_BY_RANGEandREBALANCE, which steer how many output partitions (and hence files) are produced and how evenly; with AQE turned off,REBALANCEhas no effect.