Using Adaptive Query Execution to Handle Skewed Joins in Apache Spark
Learn how Apache Spark’s Adaptive Query Execution detects skew and runtime size to switch joins, split partitions, and reduce shuffle work—complete with enablement steps, a worked example, and verification tips.
11 Aug 2026, 01:55 UTC

Problem: a join that drags because of a hot key
Imagine you are joining a large fact table (sales) with a dimension table (product) where a single product ID appears in millions of rows. The default Spark plan picks a sort‑merge join and creates many shuffle partitions, most of which are empty while a few become huge. The job runs slowly, and you see a handful of long‑running tasks in the Spark UI.
Thesis: Adaptive Query Execution (AQE) can reshape the plan at runtime to reduce shuffle work and balance the load.
AQE is enabled by default starting with Spark 3.0. After the first physical plan execution, Spark collects actual statistics (size, row count, skew) and may replan shuffles, join types, and partitioning. This can turn a sort‑merge join into a broadcast join when one side is small enough, split skewed partitions, or coalesce tiny partitions.
How to enable and verify AQE
To make sure AQE is active, set the flag and raise the log level for the adaptive planner:
# In spark-shell, pyspark, or spark-submit driver code
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.logLevel", "DEBUG")
Run the join and watch the driver logs (stdout or the log file configured by log4j.properties). You should see lines similar to:
2026-10-07 15:42:10 DEBUG AdaptiveSparkPlan: Planning adaptive query...
2026-10-07 15:42:12 INFO AdaptiveSparkPlan: Replanned stage 3 (shuffle) -> broadcast join
If you do not see these messages, AQE may be disabled or the version may be older than 3.0.
Worked example: fixing a skewed product‑sales join
Assume two Hive tables:
fact_saleswith columnssale_id, product_id, amount(≈ 200 M rows)dim_productwith columnsproduct_id, product_name, category(≈ 500 k rows, but oneproduct_idappears in 5 % of fact rows)
A naïve query:
SELECT s.amount, p.product_name
FROM fact_sales s
JOIN dim_product p ON s.product_id = p.product_id
WHERE p.category = 'Electronics'
With AQE off, Spark may choose a sort‑merge join and create 200 shuffle partitions. The hot product ends up in a single partition that processes millions of rows, while most partitions are idle.
With AQE on, the flow is:
- First execution builds a shuffle for the fact table; Spark collects runtime stats and detects that the
dim_productside after the filter is only ~2 MB. - Because this size is below the broadcast threshold (
spark.sql.autoBroadcastJoinThreshold, default 10 MB), AQE replans the join as a broadcast join. - The fact table is no longer shuffled; each executor reads the small broadcast copy of
dim_productlocally, eliminating the skew‑induced straggler.
You can compare the plans:
# With AQE enabled
spark.sql("EXPLAIN FORMATTED SELECT ...").show(truncate=false)
# With AQE disabled (for comparison)
spark.conf.set("spark.sql.adaptive.enabled", "false")
spark.sql("EXPLAIN FORMATTED SELECT ...").show(truncate=false)
The enabled version will show a BroadcastExchange node and a drastically lower number of shuffle partitions (often zero).
Trade‑offs and limitations
AQE adds a small planning overhead after the first stage. For sub‑second queries this cost can outweigh the benefit, so you may disable AQE for latency‑critical workloads:
spark.conf.set("spark.sql.adaptive.enabled", "false")Certain extensions are not fully compatible:
- Custom partitioners that rely on the original partitioning scheme may be ignored when AQE splits or coalesces partitions.
- Complex user‑defined functions (UDFs) that prevent predicate push‑down can cause AQE to fall back to the static plan.
To verify that AQE actually changed the behavior, run the same query twice—once with AQE on, once with AQE off—and compare:
- Execution time from the Spark UI (
Stagestab) or viaspark.timein Scala/PySpark. - Shuffle read/write metrics (
Shuffle ReadandShuffle Write) in the UI; a significant reduction indicates AQE helped. - Log messages for adaptive replanning as shown earlier.
Actionable closing
If you notice long‑tail tasks caused by a hot key in a join, try enabling AQE (it is on by default in recent Spark versions) and set the debug log level to confirm replanning. Compare the explain plans and UI metrics before and after. For very short jobs or when you use custom partitioners that must stay unchanged, consider turning AQE off and tuning the shuffle partition count manually.
By letting Spark adapt to the actual data distribution at runtime, you often cut shuffle volume and straggler time without rewriting the query.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.