Why does a table partitioned by dt scan every partition when a query filters only on event_ts?
answer
- the planner only sees partition columns
- two columns nobody told the table are related
- the timestamp filter becomes a row filter
- a recorded transform lets the predicate translate
basics
~20 sPruning only works on the partition column itself. A predicate on event_ts says nothing about dt, so the engine lists every partition and discards rows afterwards. Filter dt too, or use a format that derives dt from event_ts.
solid answer
~50 sThe engine prunes by evaluating predicates against **partition values**, and `dt` is just an opaque column whose relationship to `event_ts` exists only in the writer's head. Nothing in the table declares that `dt = date(event_ts)`, so a filter on `event_ts` cannot be translated into a filter on `dt`; every partition survives planning and the timestamp predicate degrades into a row filter applied after the data is read. The classic fix in a Hive-style table is to make the query carry both predicates — `WHERE dt = '2024-05-01' AND event_ts >= …` — which is fragile because every consumer must remember it. The structural fix is a format that records the transform used to derive the partition value, so the engine can convert a predicate on the source column into one on the partition value itself.
code
sql · 10 lines-- scans every partition: nothing constrains dt
SELECT count(*) FROM events
WHERE event_ts >= TIMESTAMP '2024-05-01 00:00:00'
AND event_ts < TIMESTAMP '2024-05-02 00:00:00';
-- prunes to one partition and stays exact
SELECT count(*) FROM events
WHERE dt = '2024-05-01'
AND event_ts >= TIMESTAMP '2024-05-01 00:00:00'
AND event_ts < TIMESTAMP '2024-05-02 00:00:00';go deeper
Recall that only a filter on the actual partition column skips partitions, and that a filter on a different column still returns correct results but reads everything.
Explain that pruning is a planning-time rewrite over partition predicates, and show the fix of carrying both the partition predicate and the precise timestamp predicate in the query.
Be ready to diagnose this from a query plan's partitions-read count, and to argue for moving the derivation into the table definition so every consumer prunes without special knowledge.
Own the platform framing: a layout that only performs when every query author knows a convention is an unenforced contract, and it is the reason partition derivation belongs in table metadata.
## The shape of the bug A very common lake table is partitioned by a date string `dt` that the ingestion job computes from the event's timestamp. Analysts naturally write: ```sql SELECT count(*) FROM events WHERE event_ts >= TIMESTAMP '2024-05-01 00:00:00' AND event_ts < TIMESTAMP '2024-05-02 00:00:00'; ``` The answer is correct and the query costs a full table scan. ## Why the engine cannot help Partition pruning is a rewrite performed at planning time: take the query's predicates, evaluate whichever of them reference partition columns, and drop the partitions that cannot match. The rewrite has exactly one input — predicates on partition columns. `dt` is a partition column. `event_ts` is a data column. The fact that one was computed from the other lives in the ETL code, not in the table's declaration. To the planner these are two unrelated fields, so the timestamp predicate is passed through as a residual filter applied to rows *after* they have been read, and the `dt` axis is left unconstrained. Every partition is a candidate. This is not a bug in the engine. It is the price of a partition column that is a hand-maintained duplicate of information already in the data. ## The three ways the same failure shows up **1. Filtering the source column instead of the partition column** — the case above. **2. Wrapping the partition column in a function.** `WHERE year(dt) = 2024 AND month(dt) = 5` often defeats pruning too, because many planners can only match a predicate directly against the partition value, not invert an arbitrary function over it. Some engines handle simple cases; relying on it is unwise. **3. Type and format mismatches.** If `dt` is declared `STRING` and holds `2024-05-01`, then `WHERE dt >= DATE '2024-05-01'` may force an implicit cast on the partition side, and again the planner loses its clean predicate. Worse, string ordering silently disagrees with date ordering as soon as the format is anything but zero-padded ISO. ## How to tell it is happening Read the plan. Engines report how many partitions or files survived planning versus how many exist — a "partitions read: 1461 of 1461" line, or an input-rows/bytes count equal to the whole table. If your only observation is "the query is slow", you are guessing. ## Fix one: carry both predicates ```sql WHERE dt = '2024-05-01' AND event_ts >= TIMESTAMP '2024-05-01 00:00:00' AND event_ts < TIMESTAMP '2024-05-02 00:00:00' ``` The `dt` predicate prunes; the `event_ts` predicates keep the result exact for boundary rows whose timestamp landed in a neighbouring day. This works everywhere and is what Hive-era teams learned to do reflexively. Its weakness is social: it is a rule every future query author, dashboard and generated SQL must obey, and nothing enforces it. ## Fix two: let the table own the derivation Modern table formats let the partition value be **defined as a transform of a real column** — a day-of-timestamp, a truncation, a hash bucket — and record that definition in table metadata. Two things then change: - The writer cannot get it wrong. There is no separate `dt` column for a job to compute incorrectly or forget, and no possibility of a row landing in a partition that disagrees with its own timestamp. - The planner **can** translate. A predicate on `event_ts` is pushed through the recorded transform to produce a predicate on the partition value, so the natural query prunes without the analyst knowing the layout at all. That second property is the real point: it makes the physical layout invisible to query authors, which in turn makes the layout something the table owner can change later without breaking anyone's SQL. ## The correctness footnote Even with a derived partition value, a predicate must be *monotonic-compatible* with the transform for translation to be sound — range predicates translate cleanly through day or truncate transforms, equality translates through hash buckets, but a range predicate cannot be translated through a hash bucket at all, because hashing destroys order. Expect that distinction to come up as a follow-up.
- Why is a range filter on the source column translatable through a day transform but not through a hash bucket?Day-of-timestamp is order-preserving: if a timestamp falls in a range, its day falls in the corresponding day range, so the planner can rewrite the bound. Hashing deliberately destroys order — adjacent ids land in unrelated buckets — so only equality (and IN lists) can be translated, never a range.
- The ingestion job computed dt in the local timezone while event_ts is UTC. What breaks?Rows land in a partition that disagrees with their own timestamp near midnight. A query that carries both predicates then returns too few rows, because the correct row sits in the neighbouring partition that pruning excluded. A hand-maintained partition column can be wrong; a derived one cannot.
- How do you find these queries across a whole platform rather than one at a time?Mine the engine's query history for scans whose input bytes approach the table size, or whose partitions-read equals partitions-total, and rank by total bytes. That surfaces the expensive offenders in priority order instead of relying on users to report slowness.
saying these in an interview costs you the question
- Assumes the engine infers that dt came from event_ts
- Says adding an index on event_ts would fix the pruning
- Claims the query is wrong rather than merely expensive
- Thinks wrapping the partition column in a function still prunes
- Suggests repartitioning as the only possible fix