advanced
Adaptive Query Execution (AQE)
10 min readLast updated: 2026-07-09
Overview
Learn how Adaptive Query Execution (AQE) optimizes query plans at runtime based on real-time shuffle statistics.
What You Will Learn
In this lesson, you will learn:
- Dynamic Coalescing: Reducing partition counts dynamically at runtime.
- Dynamic Join Optimization: Upgrading sort-merge joins to broadcast hash joins.
- Skew Join Mitigation: Automatically splitting skewed partitions.
Detailed Concept Explanation
Traditional databases compile query plans before execution, when exact partition sizes are unknown. Adaptive Query Execution (AQE) re-optimizes the plan during runtime, using statistics gathered during stage executions.
Key AQE Features
- Coalescing Shuffle Partitions: If the shuffle output has many small or empty partitions, AQE merges adjacent small partitions into larger ones, reducing task overhead.
- Dynamic Join Switching: If a stage reveals that a large table filtered down to under 10MB, AQE dynamically switches a Sort-Merge Join to a fast Broadcast Hash Join.
- Skew Join Mitigation: If one partition is significantly larger than the others, AQE splits it into sub-partitions, spreading the work evenly across executors and preventing straggler tasks.
Code Examples
Python (PySpark) Implementation
python
from pyspark.sql import SparkSession
# Enable AQE config
spark = SparkSession.builder \
.appName("AQETest") \
.config("spark.sql.adaptive.enabled", "true") \
.config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
.getOrCreate()
df = spark.range(1, 1000000).repartition(200)
result = df.groupBy("id").count()
# Run query to let AQE run
result.collect()
Execution Plan Diagram (Python & Scala)
Execution Plan Diagram
SparkSession.builder
config(spark.sql.adaptive.enabled
true)
range(1
1M).repartition(200)
groupBy(id).count()
collect()
Scala Implementation
scala
import org.apache.spark.sql.SparkSession
val spark = SparkSession.builder()
.appName("AQEScala")
.config("spark.sql.adaptive.enabled", "true")
.config("spark.sql.adaptive.coalescePartitions.enabled", "true")
.getOrCreate()
val df = spark.range(1, 1000000).repartition(200)
val result = df.groupBy("id").count()
result.collect()
Interview Perspective
What is Adaptive Query Execution (AQE) in Spark 3?
Adaptive Query Execution (AQE) is an optimization framework that dynamically adjusts query plans at runtime based on statistics collected during execution. Its three main features are: dynamically coalescing small shuffle partitions, dynamically switching sort-merge joins to broadcast hash joins, and dynamically splitting skewed partitions to prevent stragglers.