Apache Spark AQE: Runtime Optimization for Shuffle-Heavy Workloads
Static query planning in Apache Spark often fails because the optimizer makes decisions using stale statistics. Adaptive Query Execution (AQE) solves this by re-planning the query plan at runtime based on actual data statistics collected during shuffle operations.
04 May 2026, 08:01 UTC

Spark's query optimizer builds execution plans using table statistics that may be outdated or inaccurate. When a query contains joins or aggregations, these estimates determine the number of shuffle partitions. Too many partitions create thousands of tiny tasks, overwhelming the scheduler. Too few create straggler tasks that delay the entire job. Adaptive Query Execution (AQE) addresses this by re-optimizing the plan at runtime using actual data statistics collected during execution.
Requirements for AQE
AQE only activates after shuffle boundaries, where data is actually exchanged between executors. Local transformations like filters or projections are not re-optimized.
- Spark 3.0+: Full AQE features require Spark 3.0, with skew handling improved in 3.2+.
- Shuffle Operations: Queries must include joins, groupBy, distinct, or other shuffle-inducing operations.
- Configuration: Set
spark.sql.adaptive.enabled=trueto activate.
Three Core AQE Optimizations
AQE applies three runtime optimizations to address common distributed processing bottlenecks.
1. Dynamic Partition Coalescing
Spark's default parallelism (often 200 shuffle partitions) can create excessive task overhead for small datasets. AQE monitors actual partition sizes and merges small partitions into larger ones, reducing task count while respecting memory limits.
# Set target size for coalesced partitions (default: 64mb) spark.sql.adaptive.coalescePartitions.targetSizeInBytes=128mb
2. Join Strategy Conversion
The static optimizer may choose Sort-Merge Join based on initial table sizes. If filtering or aggregation reduces one table below a threshold, AQE can switch to Broadcast Hash Join, eliminating the expensive shuffle of the larger table.
3. Skew Join Detection
When certain keys contain disproportionate data, one executor becomes a straggler while others finish early. AQE detects skewed partitions and splits them into sub-partitions, redistributing the load across more executors.
Configuration and Trade-offs
AQE settings should align with your cluster's memory and workload characteristics. The following table shows key properties and their recommended values for skewed data workloads:
| Property | Default | Skewed Data Recommendation |
|---|---|---|
spark.sql.adaptive.enabled |
false |
true |
spark.sql.adaptive.coalescePartitions.enabled |
true |
true |
spark.sql.adaptive.skew.enabled |
false |
true |
spark.sql.adaptive.skew.joinThreshold |
2mb |
5mb |
Operational Verification
Monitor AQE effectiveness through the Spark UI's SQL tab. Look for "Adaptive Query" markers in the query plan, indicating runtime optimizations were applied. Compare execution metrics with AQE disabled to measure improvements in task duration variance and overall job time.
Limitations: AQE does not optimize local transformations. Over-coalescing can create tasks too large for executor memory if targetSizeInBytes is set too high. Skew detection may miss rapidly changing data distributions in streaming workloads.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.