Dynamic partition pruning (DPP) lets Spark skip fact table partitions using a filter that is only known at runtime, from the other side of a join. It was added in Spark 3.0 and works without you changing the query.
The problem it solves
You have sales partitioned by date_id, and a small dim_date table. A query asks for sales in Q1 2025:
SELECT s.product_id, SUM(s.amount) FROM sales s JOIN dim_date d ON s.date_id = d.date_id WHERE d.quarter = 'Q1' AND d.year = 2025 GROUP BY s.product_id;
The filter is on dim_date, not on sales. Static partition pruning cannot use it, because at planning time Spark does not know which date_id values belong to Q1. Without DPP, Spark scans every partition of sales and filters after the join.
What DPP does
At runtime, Spark first evaluates the filtered dim_date, gets the matching date_id values (about 90 of them), and passes them to the scan of sales. The scan then reads only the partitions with those dates.
dim_date --filter quarter=Q1--> {date_ids} --> sales scan reads only those partitionsWith three years of daily partitions (about 1,100), the scan reads about 90 and ignores the rest.
Conditions
- The fact table must be partitioned by the join column (here
date_id). - The join must be an equi-join on that column, and the other side must have a selective filter.
- It is on by default (
spark.sql.optimizer.dynamicPartitionPruning.enabled), and it works best with a broadcast join of the small side, though it can reuse results in other cases too.
How to confirm
In explain(), look for dynamicpruning in the PartitionFilters of the scan. In the SQL tab, the scan node shows far fewer files read than the table has.
A practical design point
This is one reason star schemas keep working well on Spark: partition the large fact table by a date key, keep dimension tables small, and let queries filter on dimension attributes like quarter or month. Make sure that date_id is the real join key and not something derived from a function, or pruning cannot apply.