Predicate pushdown means Spark pushes filter predicates from the query plan down into the data source reader, so less data is decoded and shipped to executors.
Examples
- Parquet/ORC: row-group / page-level filtering using column stats (min/max)
- JDBC sources: generate SQL
WHEREclauses executed in the database - Partition filters: skip directories (partition pruning is a form of pushdown related to layout)
Logical filter(amount > 100)
|
Physical plan pushes residual/predicate into Parquet scan
|
Reader skips non-matching row groups when stats allow# These filters can be pushed into Parquet scan
orders = spark.read.parquet("/lake/orders")
orders.filter("amount > 100 AND status = 'PAID'").select("order_id", "amount")What blocks pushdown
- Wrapping columns in opaque Python UDFs before filtering
- Some complex expressions casted in ways the source cannot understand
- Reading non-splittable / non-columnar formats (plain gzip CSV has weaker pruning)
Interview tip
"Push filters and column projections as close to the source as possible; keep data in columnar formats so pushdown and pruning actually work."