Apache Spark courseLesson 5 of 5
Apache Spark course · Lesson 5 of 5
Spark Adaptive Query Execution and Optimization
How Adaptive Query Execution re-plans Spark queries at run time: coalescing shuffle partitions, switching join strategies and splitting skewed partitions.
On this page
Spark’s optimiser plans a query before it runs, using estimates. Estimates are often wrong: a filter removes more rows than expected, or one key is far more common than others. Adaptive Query Execution (AQE) fixes this by re-optimising the remaining plan at each shuffle boundary, using actual statistics from the completed stages.
The three main features
1. Coalescing shuffle partitions
After a shuffle Spark creates spark.sql.shuffle.partitions partitions (200 by default). For small results that is far too many tiny tasks. AQE merges small adjacent partitions toward a target size (spark.sql.adaptive.advisoryPartitionSizeInBytes, 64 MB by default).
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder.master("local[2]").appName("example").getOrCreate()
counts = spark.range(0, 100_000).withColumn("k", F.col("id") % 10).groupBy("k").count()
counts.collect()
print(counts.rdd.getNumPartitions())
1
Ten tiny groups did not need 200 tasks; AQE coalesced them.
2. Switching join strategies
If, after a stage completes, one side of a join turns out to be small enough, AQE can replace a planned sort-merge join with a broadcast hash join, avoiding the shuffle of the larger side.
3. Handling skewed joins
In a sort-merge join, AQE detects partitions that are much larger than the median and splits them into smaller tasks (replicating the matching part of the other side). A partition counts as skewed when it is larger than both skewedPartitionFactor × the median (5 by default) and skewedPartitionThresholdInBytes (256 MB by default).
Settings (defaults checked on Spark 4.2)
| Setting | Default |
|---|---|
spark.sql.adaptive.enabled |
true |
spark.sql.adaptive.coalescePartitions.enabled |
true |
spark.sql.adaptive.advisoryPartitionSizeInBytes |
64 MB |
spark.sql.adaptive.skewJoin.enabled |
true |
spark.sql.adaptive.skewJoin.skewedPartitionFactor |
5 |
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes |
256 MB |
How to see AQE working
explain() before execution shows AdaptiveSparkPlan isFinalPlan=false. After the query runs, the SQL tab in the Spark UI shows the final plan, including coalesced partition counts, changed join strategies and skew-join splits.
What AQE does not fix
- Skew in aggregations (
groupByon a hot key). AQE’s skew handling targets joins; for aggregations consider salting or two-phase aggregation. - Reading too much data. Pruning and filtering early are still your job.
- Bad data models, such as joins on non-unique keys that multiply rows.
- Python UDF overhead.
Common mistakes
- Disabling AQE while debugging and forgetting to re-enable it.
- Expecting AQE to fix aggregation skew.
- Reading the pre-execution
explain()and concluding a broadcast did not happen; check the final plan.
Interview relevance
AQE is a strong follow-up to skew and shuffle questions: name its three features and its limits.
Key takeaway
AQE re-plans at shuffle boundaries using real statistics: it merges small partitions, upgrades joins to broadcast, and splits skewed join partitions. It does not replace good data layout or early filtering.
Progress is saved in this browser only. No account needed.