Architecture Note: Using Apache Spark Adaptive Query Execution in Production
Architecture note on enabling and operating Apache Spark Adaptive Query Execution: requirements, minimal config, trust boundaries, checks, failure modes, and redesign triggers.
16 Jul 2026, 03:07 UTC

Requirements
Apache Spark Adaptive Query Execution (AQE) is useful when workloads contain a mix of join types, variable data sizes, and unknown skew patterns. The goal is to let Spark re‑optimize the physical plan after map‑stage statistics are known, reducing task overhead and avoiding out‑of‑memory (OOM) failures caused by static partition sizing.
Smallest Suitable Design
The minimal configuration that activates AQE relies on three defaults introduced in Spark 3.0:
spark.sql.adaptive.enabled=true(default)spark.sql.adaptive.coalescePartitions.targetSize=134217728(128 MB)spark.sql.adaptive.autoBroadcastJoinThreshold=10485760(10 MB)spark.sql.adaptive.skewJoin.skewedPartitionFactor=5spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes=268435456(256 MB)
With these settings AQE will:
- Coalesce post‑shuffle partitions to roughly the target size.
- Convert a sort‑merge join to a broadcast hash join when the smaller side is below the broadcast threshold.
- Detect partitions that exceed
skewedPartitionFactor * median sizeand split them.
Trust / Data Boundaries
AQE only observes shuffle‑read metrics (bytes and records) from completed map tasks. It never inspects the actual column values, so personally identifiable information (PII) remains hidden. However, the observed partition sizes can leak coarse‑grained distribution patterns (e.g., that a particular key is unusually frequent). If such indirect leakage is unacceptable, AQE should be disabled.
Operational Checks
To verify that AQE is active and behaving as expected:
- In the Spark UI, open the "SQL" tab for a completed job and look for the "Adaptive Query Execution" badge.
- Run
EXPLAIN ADAPTIVEon a representative query; the output should contain anAdaptiveQueryExecutionnode and show the final number of shuffle partitions. - Enable debug logging (Spark 3.3+) with
spark.sql.adaptive.debug.enabled=trueand inspect the driver logs for lines like "Coalescing post‑shuffle partitions" or "Broadcast join selected". - Monitor executor OOM events; if they appear after enabling AQE, check whether the broadcast threshold is too large for the driver memory.
Example Configuration and Verification
Submit a Spark job with explicit AQE knobs (values can be tuned per cluster):
# Run on a gateway node with permission to submit applications
spark-submit \
--conf spark.sql.adaptive.enabled=true \
--conf spark.sql.adaptive.coalescePartitions.targetSize=134217728 \
--conf spark.sql.adaptive.autoBroadcastJoinThreshold=20971520 \
--conf spark.sql.adaptive.skewJoin.skewedPartitionFactor=4 \
--conf spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes=134217728 \
--class com.example.MyApp \
/opt/spark-apps/myapp.jar \
--input /data/events/2026-09/* \
--output /results/agg
After the job finishes, verify coalescing:
# In spark-shell or pyspark, replace with your SQL
spark.sql("EXPLAIN ADAPTIVE SELECT user_id, COUNT(*) FROM events GROUP BY user_id").show(truncate=false)
Look for a line similar to:
AdaptiveQueryExecution
+- Exchange hashpartitioning(user_id#0, 200)
The number 200 should be close to total shuffle size / targetSize. If the UI shows many more partitions than expected, increase targetSize or check for skew‑splitting that added extra tasks.
Failure Modes
- Early‑stage skew detection: In structured streaming, AQE only sees map output from the first micro‑batch; late‑arriving large partitions may be missed, causing imbalance.
- Broadcast OOM: If
autoBroadcastJoinThreshold exceeds the driver’s available memory, the broadcast side can OOM the driver. Keep the threshold below ~30% of driver memory. - Whole‑stage codegen latency: AQE‑induced plan changes disable whole‑stage codegen for the affected stages, leading to JVM recompilation spikes. Monitor GC and compilation time in the driver logs.
- Non‑deterministic functions: Queries containing
rand(),input_file_name(), or similar non‑deterministic expressions bypass AQE; the plan stays static. - Interaction with Dynamic Allocation: Frequent executor add/remove events cause AQE to repeatedly re‑coalesce partitions, leading to thrash. Stabilize executor count or disable AQE for highly elastic workloads.
Conditions That Would Change the Design
Re‑evaluate the AQE‑centric design when any of the following hold:
- The workload shifts to many short‑lived queries (< 1 second) where the overhead of collecting shuffle statistics outweighs the benefit.
- Strict latency SLAs require predictable task counts; static partitioning with manual tuning provides deterministic performance.
- You upgrade to Spark 3.4+ and want to leverage AQE’s new support for Python UDFs and additional join types; revisit thresholds because the optimizer may be more aggressive.
- Security policies prohibit any leakage of size‑based metadata; disable AQE or run with a dedicated, isolated Spark service account.
Limitations and Practical Validation
AQE cannot fix algorithmic inefficiencies (e.g., O(n²) joins) and relies on accurate map‑output size reporting. Speculative execution or task retries can duplicate metrics, causing over‑coalescing. To check for this:
- Compare
Shuffle readmetrics in the UI for a stage with and withoutspark.speculation=true. Large discrepancies suggest speculative duplication. - Run the same query with
spark.sql.adaptive.enabled=falseand compare job duration and shuffle read bytes; a significant regression indicates AQE is providing value.
If the validation shows no improvement or increased OOM events, consider disabling AQE for that workload and revert to static tuning.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.