Skip to content
Great engineers know 3 min read 1 practice problem ↓

Dynamic partition pruning

Skip whole partitions of a fact table at runtime using a filter on the dimension it joins to.

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

Comfortable with these? Read on.

TL;DR When a fact table partitioned by a key is joined to a filtered dimension, Spark uses the dimension's matching keys to skip fact partitions at runtime, even though the query never filters the fact table directly. On by default since Spark 3.0.

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:

Spark SQL
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

  1. Dimension

    Filter dates: year = 2025 and holiday

  2. Keys

    Collect the matching date_key values (say 11 keys)

  3. Prune

    Add date_key IN (those keys) as a partition filter on sales

  4. 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

DPP has nothing to prune.

Wrapping the join key in a function

A join on a cast or expression of the partition column may prevent pruning.

Expecting it without a filter on the dimension

With no selective filter there is nothing to skip.

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 questions

1. 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.

Solve: How Much Does Each Query Read? →

Go deeper

Primary sources: Configuration: dynamicPartitionPruning · Performance tuning