Menu

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.

  • Advanced
  • 16 min read
  • Updated Oct 2026
On this page
  1. Sample data
  2. Adaptive Query Execution
  3. What it is
  4. Reading AQE plans
  5. Coalescing shuffle partitions
  6. Switching from sort-merge to broadcast join
  7. Splitting skewed joins
  8. Settings (Spark 4.2 defaults)
  9. What AQE does not fix
  10. Pitfalls
  11. In interviews
  12. Dynamic partition pruning
  13. What it is
  14. When it applies
  15. Seeing it in the plan
  16. Measuring the effect
  17. Pitfalls
  18. In interviews
  19. Practice questions
  20. 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 groupBy key 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.partitions per job as if AQE did not exist. Set a sensible upper bound (or initialPartitionNum) 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_ts when the table is partitioned by order_date gets 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=true or 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 dynamicpruningexpression in PartitionFilters and the partitions-read metric.

By DataDank Editorial · Last reviewed Oct 2026 · Examples run on PySpark 4.2.0 in local mode (local[2]) with small local Parquet files. spark.sql.autoBroadcastJoinThreshold is lowered to 200 KB so a broadcast decision that needs megabytes on a cluster appears with laptop data. File paths in plans are shortened to <data>. AQE has been enabled by default since Spark 3.2.

Progress is saved in this browser only. No account needed.

Search
Filter by type