All concepts

Predicate Pushdown & Query Planning

The optimiser's whole job is to move your WHERE clause as close to the disk as it can get it.

Distributed Processing · Intermediate · ~5 min

In plain English

Telling the librarian which shelf you want before they start carrying books to your desk, instead of after.

Why it's worth your time

The difference between a 2 GB scan and a 4 TB scan is usually one clause written in a way the optimiser can use.

If you remember three things

  • Filters push down into the scan and skip partitions and row groups
  • A function around a column makes the filter opaque
  • Read the plan for PushedFilters and PartitionFilters — don't assume

Overview

You write a query as a logical statement: read this, join that, filter, group. The engine is free to execute it any way that produces the same answer, and the biggest wins come from doing less I/O. Predicate pushdown moves filters down past joins and into the scan, so rows are eliminated before they are read rather than after. Projection pushdown does the same for columns. Partition pruning skips whole directories. Together they decide whether a query reads two gigabytes or two terabytes — and the reason to understand them is that a few common query patterns silently disable them.

In an interview

Predicate pushdown moves filters down the plan — past joins, into the file scan — so the engine skips row groups and partitions instead of reading and discarding rows. Projection pushdown reads only referenced columns. Both depend on the filter being expressible against the storage layer, which is why wrapping a column in a function, or filtering on a computed alias, quietly turns a pruned scan into a full one.

Production defaults

Predicates
range comparisons on the raw column, never a function of it
Derived values
materialise frequently-filtered expressions as real columns
Review
check bytes scanned as part of code review for expensive queries

What breaks

  • Bytes scanned equals table size despite a WHERE — The predicate can't be pushed — a function, a UDF, or a filter above a window.
  • Same query, different plan after a data load — Statistics changed. That's cost-based planning working; keep stats fresh so it's working on truth.

Watch it explained

Query Optimization | SQL Query Optimization with Examples — Gate Smashers, 9:44

Related