Dynamic partition pruning
Skip whole partitions of a fact table at runtime using a filter on the dimension it joins to.
On this page
Show code in
Every code block on the page follows this.
You will learn
- What static partition pruning is, and where it stops working
- How dynamic partition pruning uses the other side of a join
- The conditions it needs
- How to confirm it in a plan
Read first
- Broadcast joins · 3 min read
- Catalyst and physical plans · 3 min read
Comfortable with these? Read on.
Static pruning, and its limit
If sales is partitioned by date_key, the query WHERE date_key = 20250301 reads one folder: that is static partition pruning, decided at planning time. But analysts rarely filter the fact table directly. They filter a dimension:
SELECT SUM(s.amount) FROM sales s JOIN dates d ON s.date_key = d.date_key WHERE d.year = 2025 AND d.is_holiday = true
Nothing in the WHERE clause mentions s.date_key, so without DPP Spark reads every partition of sales and throws most rows away in the join.
How DPP works
-
Dimension
Filter dates: year = 2025 and holiday
-
Keys
Collect the matching date_key values (say 11 keys)
-
Prune
Add date_key IN (those keys) as a partition filter on sales
-
Scan
Read only 11 partitions of sales
The pruning filter is a subquery evaluated at runtime. When the dimension is already being broadcast for the join, Spark reuses that broadcast result for pruning, so it costs almost nothing.
When it applies
- The table being pruned is partitioned by the join key (a Hive-style partition column, or a format that supports it).
- The other side has a selective filter, so pruning is worthwhile.
- Equi-join; the pruned side is not the side that must be fully preserved (for example not the left side of a left outer join).
- Enabled by
spark.sql.optimizer.dynamicPartitionPruning.enabled(true). By default Spark only applies it when it can reuse a broadcast (reuseBroadcastOnly), to avoid running the dimension query twice.
Seeing it in a plan
In the FileScan of the fact table, look for dynamicpruningexpression in PartitionFilters. In the Spark UI, the scan's "number of partitions read" metric shows the effect.
FileScan parquet sales [amount,date_key]
PartitionFilters: [isnotnull(date_key),
dynamicpruningexpression(date_key IN dynamicpruning#42)]In a lakehouse
Delta Lake and Iceberg also use file-level statistics, so a similar runtime filter can skip files, not just partitions, when the join key is clustered (Z-order, liquid clustering, sorting). The principle is the same: push the dimension's selectivity into the fact scan.
Common mistakes
Partitioning the fact table by something other than the join key
Wrapping the join key in a function
Expecting it without a filter on the dimension
Key takeaways
- DPP uses a filtered dimension to skip partitions of a fact table at runtime.
- It needs the fact table partitioned by the join key and a selective dimension filter.
- It reuses the broadcast of the dimension when possible.
- Look for dynamicpruningexpression in PartitionFilters.
Check yourself
3 questions1. What does DPP need on the fact table?
Show the answer
Partitioning by the join key. It prunes partitions, so the fact must be partitioned by the column the dimension filter maps to.
2. Where does DPP show up in explain()?
Show the answer
As dynamicpruningexpression in the fact scan PartitionFilters. The runtime filter is attached to the scan's partition filters.
3. Since which Spark version is DPP available?
Show the answer
3.0. Dynamic partition pruning was introduced in Spark 3.0.
Practice it
Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.
Go deeper
Broadcast joinsLakehouse
Partitioning done rightSpark internals
Catalyst and physical plans
Primary sources: Configuration: dynamicPartitionPruning · Performance tuning