Spark AQE Join Optimization: What Changes at Runtime and What Does Not
AQE can switch join strategies and split skewed partitions, but only at shuffle boundaries. A worked configuration, how to verify the plan changed, and where it stops helping.
17 Sept 2025, 20:52 UTC

Enabling Adaptive Query Execution (AQE) in Apache Spark does not, by itself, fix a slow join. The useful answer is narrower: AQE can change a join strategy and split skewed shuffle partitions, but only at shuffle boundaries, and only if the skew thresholds are set for your data sizes. If the slow part of your job has no shuffle before the bottleneck, AQE will not touch it.
What AQE is allowed to change
AQE is a re-optimization pass that runs between stages. Spark normally fixes a physical plan before execution starts. AQE instead waits until a shuffle map stage has written its output, reads the actual partition sizes from that shuffle, and re-plans the downstream stage. That timing is the whole constraint: no shuffle, no re-plan.
Three changes matter for joins:
- Partition coalescing — merges many small post-shuffle partitions into fewer, larger ones so tasks are not dominated by scheduling overhead.
- Join strategy switching — converts a Sort-Merge Join into a Broadcast Hash Join when one side turns out to be small after upstream filtering or aggregation.
- Skew splitting — detects a partition much larger than its peers and splits it into sub-partitions that are processed in parallel.
A worked configuration
Run this where the SparkSession is created — a notebook, a spark-submit application, or a cluster configuration file. Changing these settings requires permission to set session or cluster configuration; on a managed platform they may be fixed by the platform profile. The values below are starting points, not tuned constants.
from pyspark.sql import SparkSession
spark = (
SparkSession.builder
.appName("aqe-join-tuning")
.config("spark.sql.adaptive.enabled", "true")
.config("spark.sql.adaptive.coalescePartitions.enabled", "true")
.config("spark.sql.adaptive.skewJoin.enabled", "true")
# a partition is skewed only if it exceeds BOTH the factor and the byte threshold
.config("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5")
.config("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "256MB")
# AQE-specific broadcast ceiling (Spark 3.2+); the classic
# spark.sql.autoBroadcastJoinThreshold still applies during initial planning
.config("spark.sql.adaptive.autoBroadcastJoinThreshold", "30MB")
.getOrCreate()
)
events = spark.read.parquet("s3://bucket/events/")
users = spark.read.parquet("s3://bucket/users/")
joined = events.join(users, "user_id").filter("amount > 100")
joined.count()
Two notes on naming. The research brief behind this article referred to spark.sql.adaptive.skewJoin.threshold; in Spark 3.x the documented knob is spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes. Confirm the exact key against the configuration page for your Spark version before shipping it to a cluster, because an unrecognized key is silently ignored and you will believe a setting is active when it is not. Byte-valued settings accept strings such as 256MB; a bare integer is read as bytes, which is an easy way to set a threshold far lower than intended.
How skew splitting actually works
After the shuffle map stage completes, AQE has one size per reduce partition. A partition is treated as skewed when it is larger than both the byte threshold and the factor multiplied by the median partition size. Requiring both conditions stops AQE from splitting when every partition is uniformly large. A skewed partition is then divided into sub-partitions, and the matching rows from the other side of the join are replicated to each sub-partition so the join stays correct. Several tasks then handle one hot key instead of one task handling all of it.
Skew handling applies to the shuffle side of a Sort-Merge Join. If the join is already a Broadcast Hash Join there is no shuffle to split, and skew on the broadcast side is a different problem — usually a memory problem.
Confirming the plan changed
Do not assume the settings took effect. Three checks, in increasing order of effort:
joined.explain("formatted")shows anAdaptiveSparkPlanwrapper withisFinalPlan=false. That confirms AQE is in the plan, not that it re-planned anything useful.- The Spark UI SQL tab shows the executed plan for the query. Expand the adaptive node to see whether the final plan contains a broadcast join or split partitions. Stage-level task duration is the corroborating signal: skew splitting should flatten the long tail of task times.
- Set the log level for
org.apache.spark.sql.execution.adaptiveto DEBUG and look for re-optimization messages. This is verbose; use it in a test job rather than production.
Limits and common mistakes
- No shuffle, no AQE. Narrow transformations, a single oversized input file, and the first stage of a job are outside AQE's reach. If the bottleneck is a slow scan or a wide transformation that never shuffles, look elsewhere.
- Thresholds set too low. A small byte threshold combined with a small factor splits partitions that were never skewed, creating many short tasks and extra scheduling overhead. Raise the threshold before lowering it.
- Broadcast conversion still has a memory ceiling. AQE will not broadcast a relation larger than the relevant threshold, and broadcasting a table that fits the threshold but not executor memory can still fail the job. The adaptive broadcast threshold does not remove the memory constraint.
- Skew that never reaches a shuffle. If one input partition is huge because of how the data was written or repartitioned, AQE only sees it after it crosses a shuffle boundary.
- Assuming results are unchanged. Splitting and replication should preserve results, but verify on a sample: run the same query with
spark.sql.adaptive.enabled=falseand compare row counts and an aggregate or checksum over a fixed key range.
Rolling back
These are session or cluster settings, so rollback means restoring the previous values and restarting the session or application — there is no partial undo of an already re-planned query. Keep the prior configuration in version control so a revert is a configuration change rather than a code change.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.