Why does filtering a time-sorted columnar table by user_id prune almost no blocks?
answer
- physical order decides what you can skip
- those user rows sit everywhere in the file
- every block's id range covers the value
- min/max only helps correlated columns
- one table, one physical order
basics
~20 sBecause user ids are scattered across the whole timeline, every block's user_id min/max spans nearly the entire id domain, so no block can be proved empty. Block skipping only works on columns correlated with the table's physical order.
solid answer
~50 sBlock skipping works off per-block min/max bounds, and those bounds are only narrow for columns whose values are **correlated with the physical row order**. A table written in event-time order has each block covering a few minutes of time — a great time filter — but each of those blocks contains events from a near-random mix of users, so its `user_id` range runs from roughly the smallest id to the largest. Every block's interval covers 4711, every block survives the test, and you scan the table to return 200 rows. Nothing about the query is wrong; the *layout* is wrong for it. Fixes all change layout or add a different structure: order by `(user_id, event_time)` and give up cheap time pruning, add a per-block bloom filter for point lookups, or maintain a second copy of the table ordered by user.
code
text · 6 linesfilter: user_id = 4711
blocks scanned : 24,800 of 24,800 (no pruning)
bytes scanned : 3.9 TB
rows scanned : 12,400,000,000
rows returned : 213go deeper
Recall that skipping blocks depends on where rows physically sit: if matching rows are spread over the whole table, no block can be skipped and everything is read.
Explain correlation between the predicate column and the write order, why min/max bounds go wide when that correlation is absent, and why a lexicographic sort key's trailing columns do not fix it.
Diagnose from the profile — blocks scanned equals total, rows scanned enormous, rows returned tiny — then choose deliberately between re-sorting, a bloom filter, a partition axis, or a second physically ordered copy, and name what each costs.
Own the zero-sum framing: one table has one physical order, so every extra cheap access pattern must be bought with storage, pipeline complexity or ingest cost, and the choice should follow measured scanned bytes and query frequency.
## The shape of the problem This is the single most common pruning question in warehouse interviews because it is the single most common production surprise: a query that returns 200 rows reads 4 TB. The candidate who says "add an index" has not understood the storage model. Analytical engines skip data using per-block bounds — the min and max of each column within each block of a few hundred thousand rows. The engine tests the predicate against the interval and skips blocks it can prove cannot match. That test only helps when the interval is **narrow**, and an interval is narrow only when the column's values are **clustered** — that is, when rows with similar values were written near each other. ## Correlation with physical order is the whole story Event data typically arrives in time order and is written in time order. The consequence, per block: - `event_time` min/max: a few minutes wide out of two years. A one-day filter eliminates almost everything. - `user_id` min/max: essentially `[near-minimum id, near-maximum id]`, because any few-minute window contains events from users all over the id space. So the `user_id = 4711` predicate is compared against blocks whose ranges *all* cover 4711. No block is provably empty. The engine reads every block's `user_id` column, evaluates the filter, and discards 99.999% of the rows. This is often called an **uncorrelated predicate** — the filtered column is statistically independent of the storage order. Note the asymmetry that trips people up: the *cardinality* of `user_id` is not the problem. A high-cardinality column prunes beautifully if the table is sorted by it. What matters is scatter relative to the write order. ## Why there is no index to add Analytical stores deliberately do not maintain per-row secondary B-tree indexes: the rows are immutable, columnar, heavily compressed and often on object storage, where a million random row lookups would cost far more than one large sequential scan. There is no row pointer to chase. Everything the engine can do to avoid reading data operates at block granularity. So the levers are: **1. Change the physical order.** Order the table by `(user_id, event_time)`. Now user filters prune hard, and time filters prune well *within* a user. The costs: a full rewrite of the table (write amplification, compute time, and on some platforms a background service continuously re-sorting as new data lands), and the loss of tight time pruning for the whole-table time-range queries that every dashboard runs. Ingestion also gets more expensive, because freshly arriving rows are no longer naturally in key order and must be merged into place. **2. Keep the order, add a set-membership structure.** A per-block bloom filter (or an exact set index for lower-cardinality columns) can answer "is 4711 possibly in this block?" where min/max cannot. It rescues *equality* point lookups on a scattered column, at the cost of index size, memory and write-time work. It does nothing for range predicates and nothing for a filter that matches a large fraction of rows. **3. Combine partitioning with sorting.** Partition on the coarse axis everyone filters (day, month, region) and sort within the partition on the second axis. You get two dimensions of elimination instead of one, at the cost of more, smaller files. **4. Maintain a second physical copy.** Storage is cheap relative to repeated full scans. A derived table ordered by `user_id`, refreshed by the same pipeline, serves the per-user access pattern while the primary table stays time-ordered for analytics. The costs are pipeline complexity, freshness lag, and two things to keep correct. **5. Question the workload.** A point lookup by user id, served interactively, is an OLTP access pattern. If that is the real requirement, the honest answer may be that this data also belongs in a key-value or row store, and the warehouse should not be on that path at all. ## Diagnosing it The evidence is in the query profile: blocks or partitions scanned equals (or nearly equals) the total, bytes scanned is a large fraction of the table, and rows produced by the scan is enormous while rows returned is tiny. That gap — rows scanned versus rows output — is the fingerprint of an uncorrelated predicate. Compare it against the same query filtered on the sort column: the contrast makes the diagnosis obvious to anyone in the room. ## The generalization worth saying out loud A table has exactly one physical order, so it can have at most one predicate that prunes really well (or a few, if they are correlated with each other, like `event_time` and `event_date`). Every additional access pattern you want to serve cheaply must be bought with a different structure: a partition axis, a skipping index, a derived copy, or a pre-aggregate. Recognising that a layout is a zero-sum choice, and choosing it from measured query traffic, is the actual skill.
- Would raising or lowering the block size fix this?Only marginally, and it trades one cost for another. Smaller blocks narrow each block's ranges, so a scattered predicate skips a few more of them, but you multiply metadata, weaken compression and add per-block overhead. With a genuinely random predicate column, even tiny blocks still each contain a spread of ids. Layout, not block size, is the lever.
- The team wants to add user_id as a second sort column after event_time. Will that help?Almost not at all. Sorting is lexicographic: `user_id` is only ordered within runs sharing the same `event_time`, and with a high-precision timestamp those runs are one or two rows long. The block-level `user_id` bounds stay just as wide. Only a leading, or coarsely truncated leading, key changes block bounds.
- How do you decide between re-sorting the table and keeping a second copy ordered by user?Weigh measured traffic: bytes scanned per query times query frequency, for each access pattern. If time-range analytics dominate and user lookups are rare but latency-sensitive, keep the time order and add a derived user-ordered table or a bloom filter. If user-scoped queries dominate the bill, re-sort and accept slower whole-table time scans.
A phone book sorted by surname is useless for looking up who owns a given phone number: every page's number range covers it, so you read the whole book.
saying these in an interview costs you the question
- Blames the query optimizer instead of the physical layout
- Proposes an OLTP-style secondary index and stops there
- Thinks high cardinality alone prevents pruning
- Adds the column as a trailing sort key expecting a fix
- Ignores that re-sorting means rewriting the whole table