Apache Spark courseLesson 7 of 9
Apache Spark course · Lesson 7 of 9
Adaptive Query Execution and Dynamic Partition Pruning
How Spark re-plans queries at run time with Adaptive Query Execution, how to read AQE plans, and how dynamic partition pruning skips fact-table partitions.
On this page
- Sample data
- Adaptive Query Execution
- What it is
- Reading AQE plans
- Coalescing shuffle partitions
- Switching from sort-merge to broadcast join
- Splitting skewed joins
- Settings (Spark 4.2 defaults)
- What AQE does not fix
- Pitfalls
- In interviews
- Dynamic partition pruning
- What it is
- When it applies
- Seeing it in the plan
- Measuring the effect
- Pitfalls
- In interviews
- Practice questions
- Key takeaways
Spark’s optimiser plans a query before it runs, using estimates. Estimates are often wrong: a filter removes far more rows than expected, a join key is skewed, or a table has no statistics at all. Adaptive Query Execution (AQE) corrects this at run time by re-optimising the rest of the plan at each shuffle boundary, using the actual sizes of the stages that have finished. Dynamic partition pruning (DPP) is its companion for star-schema queries: it uses the result of a dimension filter to skip partitions of the fact table before they are read.
Sample data
A fact table of orders partitioned by date, a customer dimension of 20,000 rows (about 1.4 MB of Parquet) and a 30-row date dimension.
import contextlib, io, json, tempfile, urllib.request
from pyspark.sql import SparkSession, functions as F
spark = (SparkSession.builder.master("local[2]").appName("aqe-dpp")
.config("spark.sql.shuffle.partitions", "8")
.getOrCreate())
sc = spark.sparkContext
base = tempfile.mkdtemp()
(spark.range(0, 200_000, numPartitions=4)
.select(F.col("id").alias("order_id"),
(F.col("id") % 20000).alias("customer_id"),
((F.col("id") * 7) % 500).cast("double").alias("amount"),
F.expr("date_add(DATE'2026-01-01', cast(id % 30 as int))").alias("order_date"))
.write.partitionBy("order_date").parquet(base + "/orders"))
(spark.range(0, 20000)
.select(F.col("id").alias("customer_id"),
F.when(F.col("id") % 100 == 0, "vip").otherwise("std").alias("tier"),
F.sha2(F.col("id").cast("string"), 256).alias("pad"))
.write.parquet(base + "/customers"))
(spark.range(0, 30)
.select(F.expr("date_add(DATE'2026-01-01', cast(id as int))").alias("order_date"),
F.when(F.col("id") % 7 >= 5, "weekend").otherwise("weekday").alias("day_type"))
.write.parquet(base + "/dates"))
orders = spark.read.parquet(base + "/orders")
customers = spark.read.parquet(base + "/customers")
dates = spark.read.parquet(base + "/dates")
def show_plan(df, mode=None):
buf = io.StringIO()
with contextlib.redirect_stdout(buf):
df.explain(mode) if mode else df.explain()
print(buf.getvalue().replace("file:" + base, "<data>").rstrip())
_opener = urllib.request.build_opener(urllib.request.ProxyHandler({}))
def api(path):
port = sc.uiWebUrl.rsplit(":", 1)[1]
url = f"http://localhost:{port}/api/v1/applications/{sc.applicationId}/{path}"
return json.load(_opener.open(url))
Adaptive Query Execution
What it is
AQE, added in Spark 3.0 and enabled by default since Spark 3.2 (spark.sql.adaptive.enabled=true), splits a query into query stages at every shuffle and broadcast exchange. It runs the stages whose inputs are ready, collects their real output statistics (bytes and rows per partition), and then re-optimises and re-plans the remaining part of the query before running the next stages. Because it only acts at stage boundaries, it never changes work already done, and it can only make decisions that the shuffle statistics inform.
It has four main features:
| Feature | What triggers it | What it does |
|---|---|---|
| Coalescing shuffle partitions | Many small partitions after a shuffle | Merges adjacent small partitions towards the advisory size |
| Switching join strategy | A join side turns out small after its stage runs | Replaces a sort-merge join with a broadcast hash join (or a shuffled hash join) |
| Local shuffle reader | A join was converted to broadcast after the large side was already shuffled | Reads the already-written shuffle files locally per mapper instead of fetching across the network |
| Skew join splitting | A sort-merge join partition far larger than the median | Splits it into several tasks and replicates the matching partition of the other side |
Reading AQE plans
Before execution, explain() shows AdaptiveSparkPlan isFinalPlan=false and the initial plan. After an action, the same DataFrame’s explain() shows isFinalPlan=true with a Final Plan section, which is what actually ran. The SQL tab of the Spark UI shows the same final plan graphically. Markers to look for:
| In the final plan | Meaning |
|---|---|
ShuffleQueryStage n, BroadcastQueryStage n |
A materialised query stage, the unit AQE re-plans around |
AQEShuffleRead coalesced |
Partitions merged |
AQEShuffleRead local |
Local shuffle reader after a broadcast conversion |
AQEShuffleRead skewed, SortMergeJoin(skew=true) |
Skew join splitting applied |
BroadcastHashJoin where the initial plan had SortMergeJoin |
Join strategy switched at run time |
Coalescing shuffle partitions
After a shuffle, Spark creates spark.sql.shuffle.partitions partitions (200 by default). AQE merges contiguous small ones until each is close to spark.sql.adaptive.advisoryPartitionSizeInBytes (64 MB by default), never below spark.sql.adaptive.coalescePartitions.minPartitionSize (1 MB).
counts = orders.groupBy("order_date").count()
counts.collect()
print("partitions after the shuffle:", counts.rdd.getNumPartitions())
partitions after the shuffle: 1
Thirty small groups did not need eight tasks, so AQE ran them as one. On a cluster the same mechanism turns 200 configured partitions into however many the data actually needs. (spark.sql.adaptive.coalescePartitions.parallelismFirst, true by default, makes AQE favour parallelism and ignore the advisory size, respecting only the 1 MB minimum; the docs recommend false on busy clusters.)
Switching from sort-merge to broadcast join
The customer dimension is about 1.4 MB on disk, above the broadcast threshold (lowered to 200 KB for this demonstration), so the planner chooses a sort-merge join. Only 200 customers are vip, but without column statistics the planner cannot know how selective the filter is.
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "200KB")
vip_sales = (orders.join(customers.filter(F.col("tier") == "vip"), "customer_id")
.groupBy("tier").agg(F.sum("amount").alias("amount")))
vip_sales.collect()
show_plan(vip_sales)
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=true
+- == Final Plan ==
ResultQueryStage 4
+- *(4) HashAggregate(keys=[tier#17], functions=[sum(amount#14)])
+- AQEShuffleRead coalesced
+- ShuffleQueryStage 3
+- Exchange hashpartitioning(tier#17, 8), ENSURE_REQUIREMENTS, [plan_id=313]
+- *(3) HashAggregate(keys=[tier#17], functions=[partial_sum(amount#14)])
+- *(3) Project [amount#14, tier#17]
+- *(3) BroadcastHashJoin [customer_id#13L], [customer_id#16L], Inner, BuildRight, false, false
:- AQEShuffleRead local
: +- ShuffleQueryStage 0
: +- Exchange hashpartitioning(customer_id#13L, 8), ENSURE_REQUIREMENTS, [plan_id=183]
: +- *(1) Project [customer_id#13L, amount#14]
: +- *(1) Filter isnotnull(customer_id#13L)
: +- *(1) ColumnarToRow
: +- FileScan parquet [customer_id#13L,amount#14,order_date#15] Batched: true, DataFilters: [isnotnull(customer_id#13L)], Format: Parquet, Location: InMemoryFileIndex(1 paths)[<data>/orders], PartitionFilters: [], PushedFilters: [IsNotNull(customer_id)], ReadSchema: struct<customer_id:bigint,amount:double>
+- BroadcastQueryStage 2
+- BroadcastExchange HashedRelationBroadcastMode(List(input[0, bigint, false]),false), [plan_id=249]
+- AQEShuffleRead local
+- ShuffleQueryStage 1
+- Exchange hashpartitioning(customer_id#16L, 8), ENSURE_REQUIREMENTS, [plan_id=202]
+- *(2) Filter ((isnotnull(tier#17) AND (tier#17 = vip)) AND isnotnull(customer_id#16L))
+- *(2) ColumnarToRow
+- FileScan parquet [customer_id#16L,tier#17] Batched: true, DataFilters: [isnotnull(tier#17), (tier#17 = vip), isnotnull(customer_id#16L)], Format: Parquet, Location: InMemoryFileIndex(1 paths)[<data>/customers], PartitionFilters: [], PushedFilters: [IsNotNull(tier), EqualTo(tier,vip), IsNotNull(customer_id)], ReadSchema: struct<customer_id:bigint,tier:string>
+- == Initial Plan ==
HashAggregate(keys=[tier#17], functions=[sum(amount#14)])
+- Exchange hashpartitioning(tier#17, 8), ENSURE_REQUIREMENTS, [plan_id=148]
+- HashAggregate(keys=[tier#17], functions=[partial_sum(amount#14)])
+- Project [amount#14, tier#17]
+- SortMergeJoin [customer_id#13L], [customer_id#16L], Inner
:- Sort [customer_id#13L ASC NULLS FIRST], false, 0
: +- Exchange hashpartitioning(customer_id#13L, 8), ENSURE_REQUIREMENTS, [plan_id=140]
: +- Project [customer_id#13L, amount#14]
: +- Filter isnotnull(customer_id#13L)
: +- FileScan parquet [customer_id#13L,amount#14,order_date#15] Batched: true, DataFilters: [isnotnull(customer_id#13L)], Format: Parquet, Location: InMemoryFileIndex(1 paths)[<data>/orders], PartitionFilters: [], PushedFilters: [IsNotNull(customer_id)], ReadSchema: struct<customer_id:bigint,amount:double>
+- Sort [customer_id#16L ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(customer_id#16L, 8), ENSURE_REQUIREMENTS, [plan_id=141]
+- Filter ((isnotnull(tier#17) AND (tier#17 = vip)) AND isnotnull(customer_id#16L))
+- FileScan parquet [customer_id#16L,tier#17] Batched: true, DataFilters: [isnotnull(tier#17), (tier#17 = vip), isnotnull(customer_id#16L)], Format: Parquet, Location: InMemoryFileIndex(1 paths)[<data>/customers], PartitionFilters: [], PushedFilters: [IsNotNull(tier), EqualTo(tier,vip), IsNotNull(customer_id)], ReadSchema: struct<customer_id:bigint,tier:string>
The initial plan has SortMergeJoin with an Exchange on both sides. After the customer-side stage ran and turned out tiny, AQE re-planned the join as BroadcastHashJoin. Both sides had already been shuffled by then, so AQE used AQEShuffleRead local to read those shuffle files without a second network exchange. The larger side’s shuffle was still paid; a broadcast chosen at planning time would have avoided it, which is why accurate statistics or an explicit broadcast() hint can still beat AQE.
spark.sql.adaptive.autoBroadcastJoinThreshold lets you use a different threshold for AQE’s runtime decisions than for planning time; when unset, it falls back to spark.sql.autoBroadcastJoinThreshold.
Splitting skewed joins
A sort-merge join partition is treated as skewed when it is larger than spark.sql.adaptive.skewJoin.skewedPartitionFactor (5) times the median partition size and larger than spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes (256 MB). AQE splits it into several tasks and duplicates the matching partition on the other side. The partitions, shuffles and skew lesson demonstrates this with SortMergeJoin(skew=true) in the final plan, together with the alternatives when AQE does not apply.
Settings (Spark 4.2 defaults)
| Setting | Default | Controls |
|---|---|---|
spark.sql.adaptive.enabled |
true |
AQE as a whole |
spark.sql.adaptive.coalescePartitions.enabled |
true |
Partition coalescing |
spark.sql.adaptive.advisoryPartitionSizeInBytes |
64 MB | Target size for coalescing and skew splits |
spark.sql.adaptive.coalescePartitions.minPartitionSize |
1 MB | Smallest coalesced partition |
spark.sql.adaptive.coalescePartitions.parallelismFirst |
true |
Ignore the advisory size to keep parallelism |
spark.sql.adaptive.coalescePartitions.initialPartitionNum |
unset (uses spark.sql.shuffle.partitions) |
Starting partition count before coalescing |
spark.sql.adaptive.localShuffleReader.enabled |
true |
Local reads after broadcast conversion |
spark.sql.adaptive.skewJoin.enabled |
true |
Skew join splitting |
spark.sql.adaptive.skewJoin.skewedPartitionFactor |
5 | Skew ratio versus the median |
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes |
256 MB | Minimum size of a skewed partition |
spark.sql.adaptive.forceOptimizeSkewedJoin |
false |
Split skew even if it adds a shuffle |
spark.sql.adaptive.autoBroadcastJoinThreshold |
unset | Runtime broadcast threshold |
What AQE does not fix
- Skew in aggregations. Skew splitting targets sort-merge joins; a hot
groupBykey still lands in one task (partial aggregation usually limits the damage; otherwise salt). - Work before the first shuffle. A scan reading too much data, missing pushdown, or a bad join order among many tables is decided before any statistics exist. CBO, pruning and good layout handle that.
- The first shuffle’s cost. A join converted to broadcast at run time has already shuffled at least one side.
- Bad data models, such as joins on non-unique keys that multiply rows, and Python UDF overhead.
- Plans without exchanges. A query with no shuffle has nothing for AQE to re-plan.
Pitfalls
- Disabling AQE while debugging and forgetting to turn it back on.
- Hand-tuning
spark.sql.shuffle.partitionsper job as if AQE did not exist. Set a sensible upper bound (orinitialPartitionNum) for large jobs and let AQE coalesce. - Expecting identical job and stage counts with and without AQE: AQE runs each query stage as its own job, and the UI shows skipped stages.
- Assuming AQE applies to RDD code. It is a Spark SQL feature; RDD jobs are planned once.
In interviews
“What is AQE?” Define it (runtime re-optimisation at shuffle boundaries using real statistics), list the features (coalescing, broadcast conversion with local shuffle reader, skew join splitting), give the default status (on since 3.2), and name its limits (aggregation skew, decisions before the first shuffle). Saying that you check the final plan in the SQL tab shows practical experience.
Dynamic partition pruning
What it is
Static partition pruning skips partitions when the filter is on the partitioned table itself: WHERE order_date = '2026-01-05'. In a star schema, though, the filter is usually on a dimension: “sales on weekend days” filters dates.day_type, not orders.order_date. Without help, Spark would read every partition of orders and discard most rows in the join.
Dynamic partition pruning (added in Spark 3.0, spark.sql.optimizer.dynamicPartitionPruning.enabled, on by default) runs the dimension side first, collects the join-key values that survive its filter, and injects them as a partition filter on the fact table’s scan. Only matching partitions are listed and read.
When it applies
- The fact table is partitioned on the join key (here
order_date). - The join is an equi-join on that partition column.
- The other side has a selective filter, so pruning is worth it.
- It works best, and by default only, when the dimension side is broadcast: the broadcast result is reused as the pruning filter, so the dimension is computed once.
Seeing it in the plan
weekend_sales = orders.join(dates.filter("day_type = 'weekend'"), "order_date").agg(F.sum("amount"))
show_plan(weekend_sales)
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- HashAggregate(keys=[], functions=[sum(amount#14)])
+- Exchange SinglePartition, ENSURE_REQUIREMENTS, [plan_id=378]
+- HashAggregate(keys=[], functions=[partial_sum(amount#14)])
+- Project [amount#14]
+- BroadcastHashJoin [order_date#15], [order_date#19], Inner, BuildRight, false, false
:- FileScan parquet [amount#14,order_date#15] Batched: true, DataFilters: [], Format: Parquet, Location: InMemoryFileIndex(1 paths)[<data>/orders], PartitionFilters: [isnotnull(order_date#15), dynamicpruningexpression(order_date#15 IN dynamicpruning#53)], PushedFilters: [], ReadSchema: struct<amount:double>
: +- SubqueryAdaptiveBroadcast dynamicpruning#53, [0], true, Project [order_date#19], [order_date#19]
: +- AdaptiveSparkPlan isFinalPlan=false
: +- Project [order_date#19]
: +- Filter ((isnotnull(day_type#20) AND (day_type#20 = weekend)) AND isnotnull(order_date#19))
: +- FileScan parquet [order_date#19,day_type#20] Batched: true, DataFilters: [isnotnull(day_type#20), (day_type#20 = weekend), isnotnull(order_date#19)], Format: Parquet, Location: InMemoryFileIndex(1 paths)[<data>/dates], PartitionFilters: [], PushedFilters: [IsNotNull(day_type), EqualTo(day_type,weekend), IsNotNull(order_date)], ReadSchema: struct<order_date:date,day_type:string>
+- BroadcastExchange HashedRelationBroadcastMode(List(input[0, date, true]),false), [plan_id=373]
+- Project [order_date#19]
+- Filter ((isnotnull(day_type#20) AND (day_type#20 = weekend)) AND isnotnull(order_date#19))
+- FileScan parquet [order_date#19,day_type#20] Batched: true, DataFilters: [isnotnull(day_type#20), (day_type#20 = weekend), isnotnull(order_date#19)], Format: Parquet, Location: InMemoryFileIndex(1 paths)[<data>/dates], PartitionFilters: [], PushedFilters: [IsNotNull(day_type), EqualTo(day_type,weekend), IsNotNull(order_date)], ReadSchema: struct<order_date:date,day_type:string>
The fact scan’s PartitionFilters contains dynamicpruningexpression(order_date IN dynamicpruning#...), fed by SubqueryAdaptiveBroadcast, which reuses the dimension’s broadcast exchange. In explain("formatted") output the subquery appears under ===== Subqueries ===== with a ReusedExchange.
Measuring the effect
The SQL tab shows scan metrics for each query; the same numbers come from the REST API. Running the query with DPP on and off:
def fact_scan_metrics(description):
for q in api("sql?details=true"):
if q["description"] == description:
for node in q["nodes"]:
metrics = {m["name"]: m["value"] for m in node["metrics"]}
if node["nodeName"].startswith("Scan parquet") and "number of partitions read" in metrics:
return {k: metrics[k] for k in
["number of partitions read", "number of files read", "number of output rows"]}
for enabled in ["true", "false"]:
spark.conf.set("spark.sql.optimizer.dynamicPartitionPruning.enabled", enabled)
sc.setJobDescription("DPP " + enabled)
result = orders.join(dates.filter("day_type = 'weekend'"), "order_date").agg(F.sum("amount")).collect()
sc.setJobDescription(None)
print("DPP", enabled, result[0][0], fact_scan_metrics("DPP " + enabled))
spark.conf.set("spark.sql.optimizer.dynamicPartitionPruning.enabled", "true")
DPP true 13240021.0 {'number of partitions read': '8', 'number of files read': '32', 'number of output rows': '53,333'}
DPP false 13240021.0 {'number of partitions read': '30', 'number of files read': '120', 'number of output rows': '200,000'}
Same answer, but with DPP the fact scan read 8 of 30 date partitions (the weekend days) and about a quarter of the rows. On a fact table of terabytes, that is the difference between minutes and hours.
Pitfalls
- The fact table is not partitioned by the join key. DPP prunes directories, not row groups. Joining on
order_tswhen the table is partitioned byorder_dategets no pruning; join or filter on the partition column. - No selective filter on the dimension. Joining an unfiltered dimension prunes nothing.
- Dimension too big to broadcast. By default DPP only uses a broadcast it can reuse; a large dimension may get no pruning. Filter or project it down so it broadcasts.
- Functions on the partition column in the join condition (
to_date(o.ts) = d.date) can stop the rewrite. - Confusing it with runtime filters for non-partitioned tables. Spark also has runtime Bloom filters (
spark.sql.optimizer.runtime.bloomFilter.enabled, on in 4.2) that can reduce rows from the large side of a shuffle join; DPP specifically prunes partitions.
In interviews
“What is dynamic partition pruning?” Explain the star-schema problem (filter on the dimension, partitioning on the fact), the mechanism (evaluate the dimension filter, push the surviving keys into the fact scan as a partition filter, reusing the broadcast), the requirements (partitioned fact table, equi-join on the partition column, selective dimension filter), and how to prove it worked (dynamicpruningexpression in PartitionFilters, fewer partitions read in the SQL tab).
Practice questions
At what points does AQE re-optimise, and with what information?
At query-stage boundaries, which are shuffle and broadcast exchanges. When a stage finishes, AQE has its actual output statistics (bytes and rows per partition), re-runs optimisation on the remaining logical plan with those statistics, and re-plans the physical operators before submitting the next stages. Completed work is never redone.
The pre-execution plan shows SortMergeJoin, but the job ran faster than expected. How do you find out what happened?
Look at the final plan: call explain() on the same DataFrame after the action (it shows isFinalPlan=true), or open the query in the SQL tab. If AQE converted the join, the final plan shows BroadcastHashJoin, a BroadcastQueryStage, and usually AQEShuffleRead local on the side that had already been shuffled.
Why is a broadcast join chosen at planning time cheaper than one chosen by AQE at run time?
AQE only learns that a side is small after that side’s shuffle stage has run, and by then the large side may have been shuffled too. The local shuffle reader avoids a second network fetch, but the first shuffle was already paid. A planning-time broadcast (from statistics or a broadcast() hint) never shuffles the large side.
Which problems does AQE not solve?
Skewed aggregations (skew splitting is for sort-merge joins), anything decided before the first shuffle (excess scanning, missing pushdown, multi-way join order), the cost of the shuffle that produced the statistics, data-model problems such as fan-out joins, Python UDF overhead, and RDD jobs.
A query joins a sales table partitioned by sale_date to a calendar filtered to public holidays, but it reads every partition. What do you check?
That the join is an equi-join directly on sale_date (not on a timestamp or a function of it); that DPP is enabled; that the filtered calendar is small enough to broadcast so its result can be reused for pruning; and that the calendar side really has the filter. Then confirm dynamicpruningexpression appears in the sales scan’s PartitionFilters and that the SQL tab shows fewer partitions read.
What is the difference between static and dynamic partition pruning?
Static pruning uses a literal filter on the partition column of the scanned table, known at planning time. Dynamic pruning derives the partition filter at run time from another table’s filtered join keys, so a filter on a dimension can prune the fact table.
Key takeaways
- AQE re-plans at each shuffle or broadcast boundary using measured statistics; it has been on by default since Spark 3.2.
- Its features are partition coalescing, runtime conversion to broadcast joins (with a local shuffle reader), and skew join splitting.
- Always read the final plan (
isFinalPlan=trueor the SQL tab) to see what AQE did. - AQE does not fix aggregation skew, scanning too much data, or decisions made before the first shuffle.
- Dynamic partition pruning turns a dimension filter into a partition filter on the fact table; it needs a fact table partitioned on the join key and, by default, a broadcastable dimension.
- Confirm DPP with
dynamicpruningexpressioninPartitionFiltersand the partitions-read metric.
Progress is saved in this browser only. No account needed.

