What is late materialization in a columnar engine, and when does it beat fetching projected columns up front?
answer
- the order of filter and fetch is a choice
- a position list is cheaper than a tuple
- the payoff needs two things, not one
- if most rows survive you pay twice
- encoded chunks resist jumping to position k
basics
~20 sLate materialization evaluates predicates on the filter columns first, producing a list of surviving row positions, and only then fetches the projected columns at those positions. It wins when the filter is selective and the projected columns are wide.
solid answer
~50 s**Early materialization** reads every projected column for every candidate row, assembles tuples, then applies the predicate — so work is done for rows that are immediately discarded. **Late materialization** reverses the order: it decodes only the columns the predicate needs, evaluates the filter to produce a list of surviving positions within each row group, and gathers the projected columns only at those positions. The payoff scales with two things: how few rows survive, and how expensive the deferred columns are to read and decode. Filtering to 0.5% of rows on a table whose payload column is a 2 KB blob means avoiding roughly 99.5% of that blob's decoding. It stops paying when selectivity is poor — if most rows survive, you have made two passes instead of one and gained nothing — and when the deferred column's encoding does not support cheap positional access, so a scattered gather ends up decoding the chunk anyway.
code
sql · 7 lines-- cheap narrow filter, one very expensive projected column
SELECT session_id, request_payload -- request_payload ~2 KB per row
FROM events
WHERE status_code = 500; -- keeps roughly 0.3% of rows
-- deferring request_payload until the surviving positions are known
-- avoids decoding it for the ~99.7% of rows that are discardedgo deeper
Know the distinction in plain terms: fetch the wide columns before filtering, or after. Filtering first and fetching only the surviving rows is usually the cheaper order.
Explain the mechanism — a position list rather than tuples, produced by evaluating predicates on the filter columns first — and state that the payoff depends on both selectivity and the width of the deferred columns.
Show the failure modes: high survivor rates turning it into two passes, encodings that defeat positional access, scattered gathers becoming request-bound on remote storage, and how physical ordering flips the outcome.
Own the point that this hinges on optimizer selectivity estimates, which are unreliable on correlated analytical predicates — so table layout and projection discipline are more dependable levers than trusting the engine to pick right.
## The two orders of operations A scan of a columnar table has to do three things: read the columns the predicate needs, evaluate the predicate, and produce the projected columns for the rows that survive. The only real freedom is *when* the projected columns get read. **Early materialization** reads all referenced columns for the candidate rows, stitches them into tuples by ordinal position, and then filters the tuples. It is simple, it produces one sequential pass over every chunk, and every operator downstream sees ordinary rows. **Late materialization** keeps the columns apart as long as possible. It reads only the predicate columns, evaluates the filter against them, and produces not tuples but a **list of surviving positions** within the row group. The projected columns are fetched afterwards, at those positions only, and stitching happens at the last possible moment — sometimes after aggregation has already consumed the values, in which case some columns are never assembled into a row at all. ## Why deferring pays The saving is the product of two factors. **Selectivity.** If a predicate on a narrow, cheap column keeps 1 row in 500, then 499 of every 500 reads and decodes of the projected columns were avoidable work, and late materialization avoids them. **The cost of the deferred columns.** Columns are not equal. A filter on an integer date column is cheap; a projected JSON, free-text or wide binary column may be orders of magnitude more expensive per row to fetch and decode. Deferring the expensive column behind the cheap filter is the whole trick, and it is why the technique matters most on tables that mix a few narrow filter columns with a few very wide payload columns. There is a second, related benefit: while values stay as positions rather than tuples, intermediate results are compact. A surviving-position list for a slice is far smaller than the same rows materialized, so passing it between operators costs less memory and less copying. ## Why it is not always chosen **Poor selectivity turns it into a loss.** If 80% of rows survive, the position list is nearly as long as the slice and you have replaced one sequential pass with a predicate pass plus a near-full gather. Two passes over data you were going to read anyway is strictly worse than one, and the gather is less cache-friendly than the sequential read. **Encoding may not support positional access.** Reaching position *k* is trivial in a fixed-width uncompressed chunk. In a run-length, delta or block-compressed chunk it may require decoding from the start of the enclosing block. A scattered gather across many blocks can therefore decode most of the chunk regardless, and the saving evaporates. Engines mitigate this by grouping positions per block and by preferring the technique on columns whose layout supports skipping. **Random access patterns.** On remote storage, a gather at scattered positions can turn one large sequential read into many small ranged reads, where per-request latency dominates. Engines commonly coalesce nearby positions back into ranges, which reintroduces some wasted bytes but restores throughput — a middle ground between the two extremes. **Multiple predicates change the calculus.** With several filter columns, engines typically evaluate the cheapest and most selective predicate first, narrow the position list, then evaluate the next predicate only at the surviving positions. That is late materialization applied within the filter set, and it is where much of the win often lives — the second predicate column is itself a deferred column. ## How the optimizer decides The decision rests on estimated selectivity, per-column width, and whether the storage layer supports cheap positional skipping. Because selectivity estimates on analytical data are frequently wrong — correlated predicates, skewed values — some engines decide adaptively at runtime: evaluate the predicate on the first portion of a slice, measure the actual survival rate, and pick the strategy for the rest. A candidate who mentions that the decision is estimate-driven and therefore fallible is answering at the right level. ## Diagnosing it Symptoms that a scan is materializing too early: the query reads far more bytes than the surviving row count suggests it should; a query slows dramatically when a wide column is added to the SELECT list even though the filter and row count are unchanged; profiles show decode time dominated by a column that contributes nothing to the filter. The practical mitigations do not all require the engine to be clever: dropping an unused wide column from the projection removes the problem entirely, and physically ordering the data so the predicate's matching rows cluster together turns a scattered gather back into a sequential one. ## Interview framing State the definition in one sentence, name the two factors that determine the payoff (selectivity and deferred-column cost), and then — this is what separates a strong answer — name the two regimes where it loses: a high survivor rate, and encodings that make positional access as expensive as sequential decoding. Finish with the observation that the technique is only worth reasoning about on tables where a cheap filter column guards an expensive payload column; on a table of ten integers it does not matter.
- Why does the benefit collapse when the filter keeps 80% of rows?Because you have replaced a single sequential pass with a predicate pass plus a gather that touches almost everything anyway. The position list is nearly as long as the slice, the gather is less cache- and prefetch-friendly than a straight walk, and almost none of the deferred decoding was actually avoided. Early materialization is the cheaper plan in that regime.
- How does compression interact with fetching values at scattered positions?Badly, in general. Block-compressed or run-length-encoded chunks do not support arithmetic indexing: reaching position k means decoding from the start of its block. A gather spread thinly across many blocks therefore decodes most of the chunk, erasing the saving. Engines reduce the damage by sorting positions and coalescing adjacent ones into ranged reads.
- Does clustering the data on the filter column change the picture?Substantially. If matching rows are physically adjacent, the surviving positions form contiguous runs rather than scattered singletons, so the deferred fetch becomes a small number of sequential range reads instead of many random ones. Physical ordering converts the worst case of late materialization into its best case.
- Can materialization be deferred past a join?Some engines do it: the join runs on key columns and produces position or row-identifier pairs, and the wide payload columns are fetched afterwards only for matched rows. It pays for the same reason as in a scan — the join is often far more selective than its inputs — but it requires carrying positions through the join operator, which not all engines support.
saying these in an interview costs you the question
- Says late materialization is always faster than early
- Ignores that a poorly selective filter makes it a double pass
- Assumes any position in an encoded chunk can be read directly
- Confuses it with projection pushdown or predicate pushdown
- Claims it helps even when every projected column is a narrow integer