Spark Adaptive Query Execution: Deep Dive
Apache Spark's query optimizer has always been impressive, but it has a fundamental limitation: it makes all optimization decisions before a single row of data is processed. The Catalyst optimizer builds a physical plan based on table statistics, heuristics, and configuration parameters. If those statistics are stale, missing, or misleading, the resulting plan can be catastrophically suboptimal. Adaptive Query Execution (AQE) changes this by re-optimizing the query plan at runtime, using actual data statistics collected at stage boundaries.
Enabled by default since Spark 3.2, AQE addresses three critical problems that have plagued Spark performance tuning for years: over-provisioned shuffle partitions, data skew in joins, and missed broadcast join opportunities. This article explores each mechanism in detail, with practical configurations and code examples drawn from production workloads.
The Static Optimization Problem
Before AQE, tuning a Spark job was an exercise in educated guessing. The default value of spark.sql.shuffle.partitions is 200, which is too few for terabyte-scale datasets and too many for gigabyte-scale ones. Engineers would profile jobs, examine shuffle read metrics, and manually adjust partition counts. This process had to be repeated whenever data volumes changed significantly.
The Catalyst optimizer relies on statistics collected by ANALYZE TABLE commands or inferred from file metadata. In a typical data warehouse environment, these statistics become stale within hours as new data arrives. Consider a pipeline that joins a slowly-changing dimension table against a rapidly-growing fact table:
-- Static optimizer sees the dimension table as 500MB
-- based on statistics collected last week.
-- Today it has grown to 4GB after a schema migration
-- added denormalized columns.
SELECT f.*, d.category_name, d.region
FROM fact_events f
JOIN dim_products d ON f.product_id = d.product_id
WHERE f.event_date = '2026-09-30'
With stale statistics, the optimizer might choose a broadcast hash join for the dimension table, attempting to broadcast 4GB to every executor. The result is either an out-of-memory failure or extreme garbage collection pressure. AQE solves this by checking actual data sizes at runtime and switching join strategies when the original choice no longer makes sense.
Dynamic Partition Coalescing
The most immediately impactful AQE feature is dynamic coalescing of shuffle partitions. Instead of guessing the right number of partitions, you set a deliberately high initial count and let AQE merge small partitions together after the shuffle completes.
The key configuration parameters are:
# Enable AQE (default: true since Spark 3.2)
spark.sql.adaptive.enabled = true
# Enable partition coalescing
spark.sql.adaptive.coalescePartitions.enabled = true
# Target size for coalesced partitions (default: 64MB)
spark.sql.adaptive.advisoryPartitionSizeInBytes = 128MB
# Initial shuffle partitions (set high, AQE coalesces down)
spark.sql.shuffle.partitions = 2000
# Minimum number of partitions after coalescing
spark.sql.adaptive.coalescePartitions.minPartitionNum = 1
The algorithm works by scanning the shuffle partition sizes from the completed map stage. It greedily combines adjacent partitions until the combined size approaches the advisory target. This means a job processing 10GB of shuffle data with a 128MB target will end up with roughly 80 partitions, regardless of the initial setting.
In PySpark, you can observe the effect by examining the physical plan:
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.config("spark.sql.adaptive.enabled", "true") \
.config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
.config("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128MB") \
.config("spark.sql.shuffle.partitions", "2000") \
.getOrCreate()
df = spark.read.parquet("s3://warehouse/fact_events/")
result = df.groupBy("product_id", "region") \
.agg({"revenue": "sum", "event_id": "count"})
# Force execution and examine the adaptive plan
result.write.mode("overwrite").parquet("s3://warehouse/agg_output/")
# Check the Spark UI: the "AQE" tab shows
# the original vs. optimized partition counts
The performance impact is substantial. A common pattern in production is a pipeline with highly variable data volumes across partitions, for example when grouping by customer ID where a few large enterprise accounts dominate the dataset. Without coalescing, you end up with thousands of tiny partitions that each incur scheduler overhead. With coalescing, those partitions merge into well-sized chunks that execute efficiently.
Skew Join Optimization
Data skew is the single most common cause of Spark job straggler tasks. A join where one key appears millions of times while most keys appear once creates a partition that takes orders of magnitude longer than the others. Before AQE, the standard remedies were salting (adding random suffixes to keys) or isolating skewed keys into separate joins.
AQE's skew join optimization detects skew automatically. After the map stage completes, it compares each partition's size against the median. If a partition exceeds a configurable threshold, AQE splits it into smaller sub-partitions and replicates the corresponding partition from the other side of the join.
# Skew join configuration
spark.sql.adaptive.skewJoin.enabled = true
# A partition is skewed if it is larger than this
# multiple of the median partition size
spark.sql.adaptive.skewJoin.skewedPartitionFactor = 5
# AND larger than this absolute threshold
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes = 256MB
Both conditions must be true for a partition to be classified as skewed. This dual-threshold approach prevents false positives: a uniformly large dataset where all partitions are 500MB will not trigger skew handling, even though each partition exceeds the absolute threshold, because no partition is significantly larger than the median.
Here is a practical example. Suppose you are joining user events against a user profiles table, and a handful of bot accounts generate millions of events:
-- Without AQE, the partition for bot user IDs
-- might contain 50 million rows while normal partitions
-- have 50,000 rows each.
SELECT e.event_type,
u.account_tier,
COUNT(*) as event_count,
AVG(e.duration_ms) as avg_duration
FROM events e
JOIN user_profiles u ON e.user_id = u.user_id
WHERE e.event_date BETWEEN '2026-09-01' AND '2026-09-30'
GROUP BY e.event_type, u.account_tier
With AQE skew join enabled, the oversized partitions for bot user IDs are automatically split. The corresponding user_profiles partitions are replicated to match. The result is even task distribution without any changes to the query logic. This is particularly valuable in environments where skew patterns shift over time, making manual salting strategies brittle.
Runtime Broadcast Join Conversion
Broadcast hash joins are the fastest join strategy in Spark: one side of the join is sent to every executor, eliminating the shuffle entirely. The static optimizer chooses a broadcast join when it estimates one side is smaller than spark.sql.autoBroadcastJoinThreshold (default 10MB). But size estimates from statistics can be wildly inaccurate, especially after complex filter predicates.
AQE re-evaluates join strategies after the map stage, using the actual data size rather than the estimated size. If one side of a sort-merge join turns out to be small enough for a broadcast, AQE converts it. This frequently happens with filtered dimension tables:
-- The dim_stores table is 2GB total (above broadcast threshold)
-- But filtering to a single region yields only 5MB
SELECT f.transaction_id,
f.amount,
s.store_name,
s.city
FROM fact_transactions f
JOIN dim_stores s ON f.store_id = s.store_id
WHERE s.region = 'APAC'
AND f.transaction_date = '2026-09-30'
The static optimizer sees dim_stores as 2GB and plans a sort-merge join. But after the filter pushdown executes, only 5MB of data remains. AQE detects this and converts to a broadcast hash join, eliminating the shuffle of the fact table entirely. On a 500-node cluster processing terabytes of transactions, this conversion can reduce query time from minutes to seconds.
When working with high-performance data transfer protocols like Arrow Flight, the ability to dynamically select join strategies becomes even more important, as data arrives with unpredictable volume patterns.
AQE and the Query Stage Model
Understanding AQE requires understanding Spark's stage execution model. A Spark job is divided into stages at shuffle boundaries. Each stage consists of a set of tasks that execute the same computation on different partitions. AQE operates at stage boundaries because that is where Spark has complete information about the intermediate results.
The AQE optimization loop works as follows:
- The Catalyst optimizer generates an initial physical plan as usual.
- AQE wraps each query stage in a
QueryStageExecnode that can be re-optimized. - Leaf stages (those with no shuffle dependencies) execute first.
- When a stage completes, AQE collects accurate statistics: partition sizes, row counts, and data distribution.
- Using these runtime statistics, AQE re-runs optimization rules on the remaining unexecuted stages.
- The updated plan may change join strategies, coalesce partitions, or split skewed partitions.
- The next set of stages executes with the optimized plan.
This stage-by-stage approach means AQE's benefits compound in multi-join queries. Consider a query that joins five tables in sequence. After the first join completes, AQE knows the exact output size and can optimize the second join accordingly. Each subsequent join benefits from increasingly accurate statistics.
AQE is most impactful on complex queries with multiple joins and aggregations. A simple single-table scan and filter gets minimal benefit because there are no shuffle stages where re-optimization can occur.
Production Tuning Strategies
While AQE's defaults work well for many workloads, production environments benefit from deliberate tuning. Here is a configuration template that works well for large-scale data warehouse workloads:
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("production_etl") \
.config("spark.sql.adaptive.enabled", "true") \
.config("spark.sql.adaptive.forceOptimizeSkewedJoin", "true") \
.config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
.config("spark.sql.adaptive.coalescePartitions.parallelismFirst", "false") \
.config("spark.sql.adaptive.advisoryPartitionSizeInBytes", "256MB") \
.config("spark.sql.adaptive.skewJoin.enabled", "true") \
.config("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5") \
.config("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "256MB") \
.config("spark.sql.adaptive.localShuffleReader.enabled", "true") \
.config("spark.sql.shuffle.partitions", "4000") \
.config("spark.sql.autoBroadcastJoinThreshold", "64MB") \
.getOrCreate()
Several of these settings deserve explanation:
- parallelismFirst = false: By default, AQE prioritizes parallelism over partition size. Setting this to false tells AQE to prioritize reaching the advisory partition size, which typically yields better throughput for I/O-heavy ETL workloads.
- advisoryPartitionSizeInBytes = 256MB: For large datasets, bigger partitions reduce scheduler overhead. The default 64MB is conservative; 128-256MB works better on modern executors with 8+ GB memory.
- shuffle.partitions = 4000: Set this higher than you expect to need. AQE will coalesce down. The cost of starting high is minimal compared to the risk of too few partitions for the occasional large dataset.
- localShuffleReader.enabled = true: Allows AQE to convert remote shuffle reads to local reads when the data is co-located, reducing network traffic after broadcast join conversion.
For pipelines that feed Delta Lake change data feeds, proper AQE tuning ensures that the write stage produces appropriately-sized files, which directly impacts downstream read performance.
Monitoring and Debugging AQE
AQE transformations are visible in the Spark UI under the SQL tab. Each query shows both the initial plan and the final adaptive plan. Key metrics to monitor include:
| Metric | What It Shows | Action If Abnormal |
|---|---|---|
| Coalesced partitions | How many original partitions were merged | If no coalescing occurs, initial partition count may already be optimal |
| Skewed partitions detected | Number of partitions that exceeded skew thresholds | High counts suggest persistent skew; consider pre-aggregation |
| Join strategy changes | Conversions from sort-merge to broadcast | Frequent conversions indicate stale table statistics |
| Stage replanning time | Overhead of AQE re-optimization | Usually under 100ms; if higher, check for overly complex plans |
You can also inspect AQE decisions programmatically:
# After executing a query, examine the adaptive plan
df = spark.sql("""
SELECT product_id, SUM(revenue) as total_revenue
FROM fact_events
GROUP BY product_id
HAVING SUM(revenue) > 10000
""")
# Trigger execution
df.collect()
# The explain output shows AQE annotations
df.explain("formatted")
# Look for:
# AdaptiveSparkPlan isFinalPlan=true
# CustomShuffleReader coalesced
# SkewedJoin markers
When AQE is not behaving as expected, the most common culprit is the spark.sql.adaptive.coalescePartitions.parallelismFirst setting. When set to true (the default), AQE preserves the Spark default parallelism as a floor and will not coalesce below that number of partitions, even if the data would fit in fewer.
For teams building real-time analytics on ClickHouse, understanding how AQE optimizes the upstream Spark processing that feeds those materialized views is essential for end-to-end pipeline efficiency.
Limitations and Considerations
AQE is not a silver bullet. Several important limitations affect production usage:
- Stage boundary requirement: AQE can only re-optimize at shuffle boundaries. Operations within a single stage, such as a chain of map transformations, cannot be re-optimized.
- No cross-job learning: AQE does not persist runtime statistics between job executions. Each run starts fresh, making the same discoveries again. For recurring jobs, combining AQE with periodically refreshed table statistics yields the best results.
- Streaming limitations: AQE works with Structured Streaming in micro-batch mode but not in continuous processing mode. The micro-batch boundaries serve as the stage boundaries where re-optimization occurs.
- Non-deterministic plans: Because AQE adapts to runtime conditions, the same query on the same data might produce different physical plans if executed at different cluster loads. This complicates debugging but rarely affects correctness.
- Overhead on small queries: For queries that complete in seconds, the AQE overhead of collecting statistics and re-optimizing can be a noticeable percentage of total execution time. Consider disabling AQE for interactive, sub-second dashboarding queries.
Despite these limitations, AQE has become the single most impactful optimization feature in modern Spark deployments. It replaces hours of manual tuning with automatic, runtime-aware decisions that adapt to changing data characteristics. For any team running Spark at scale, understanding and properly configuring AQE is essential knowledge.