A columnar file's footer records each column's minimum and maximum per chunk of rows; how can a reader use that to leave most of the file unread?
answer
- metadata answers before the data is fetched
- footer records per-chunk value ranges
- no overlap means the chunk is skipped
- overlap only proves a match is possible
- write order decides how much is skipped
basics
~20 sThe reader evaluates the query's filter against each chunk's recorded value range first. A chunk whose range cannot intersect the filter is skipped without being fetched at all; only chunks whose range overlaps are read, and their rows are then re-checked individually.
solid answer
~50 sA columnar file is written in chunks of rows, and its footer describes every chunk: where each column's block sits and what that block contains — minimum, maximum, null count. Because the footer is small and read first, the reader can ask a cheap question per chunk before fetching anything: could any row here satisfy the filter? Three outcomes follow. The chunk's range lies entirely outside the filter, so it is skipped and never fetched. The range lies entirely inside, so every row qualifies and no per-row comparison is needed. The ranges overlap, so the chunk must be read and its rows tested one by one. How much this saves is decided not by the format but by the data's order: if the filtered column is scattered across the write order, every chunk's range spans the whole domain and nothing can be skipped.
code
pseudocode · 13 linesmatched = 0
for each chunk in footer.chunks:
stats = chunk.stats_for("event_time")
if stats is missing:
read_and_filter(chunk) # no statistics: never skip
continue
if stats.max < range_start or stats.min > range_end:
continue # cannot match: blocks never fetched
times = read(chunk, "event_time") # may match: read and re-check
amounts = read(chunk, "amount")
for i in 0 .. times.length - 1:
if times[i] >= range_start and times[i] <= range_end:
matched = matched + amounts[i]go deeper
Recall that a columnar file ends with a footer describing every chunk of rows, including the smallest and largest value each column holds there, and that a reader consults it before fetching any data.
Explain the three-way outcome of comparing a chunk's range against a filter — disjoint, contained, overlapping — and why an overlap still requires testing each row individually.
Show the diagnosis. Explain why a filtered query can still read everything when the data is clustered by a different column, and what evidence you would gather before blaming the format.
Own the clustering policy. Decide which column a dataset is ordered by given the query mix, what chunk size trades away, and whether a second differently ordered copy is worth its storage and synchronisation cost.
## The shape the footer describes A columnar file is not one enormous column per file. The writer accumulates a **chunk** of rows — typically hundreds of thousands to a few million — and, when the chunk is full, writes one contiguous block per column for exactly those rows, then starts the next chunk. The file therefore reads as a sequence of chunks, each internally column-major. The **footer** at the end of the file is the map: the schema, the byte offset and length of every column's block in every chunk, and a small statistics record per column per chunk. Those statistics are what make skipping possible. They typically carry: - the **minimum** and **maximum** value present in that block, - the **null count**, which answers `IS NULL` and `IS NOT NULL` style filters directly, - sometimes a **distinct-value count** or a small sketch, where the writer chose to compute one. A reader opens the file by reading the footer, which costs one small request, and only then decides what else to fetch. ## The three-way decision per chunk For a filter on one column, each chunk falls into one of three cases: 1. **Disjoint.** The chunk's `[min, max]` does not intersect the filter's range, so no row in it can qualify. The chunk's blocks are never requested. This is where the saving comes from. 2. **Contained.** The chunk's `[min, max]` lies entirely inside the filter's range, so every row qualifies. The column still has to be read if it is projected, but the per-row comparison can be skipped. 3. **Overlapping.** The ranges intersect. This proves *nothing* about any individual row: the chunk must be read and every row tested. Treating an overlap as a match is a real and costly error. A fourth case matters in practice: **statistics absent or untrusted**. A writer may omit them, and value ranges for some types are not always comparable in the way a reader expects. The safe default is always to read the chunk. A reader that treats a missing statistic as permission to skip silently returns wrong results, which is far worse than being slow. Worked numbers: one billion rows written in timestamp order, chunks of one million rows, so a thousand chunks. A year of data means roughly 2.74 million rows per day, so a one-day filter intersects about three consecutive chunks. The reader fetches three chunks out of a thousand — roughly 0.3% of the file — plus the footer. ## Why order, not the format, decides the saving Now write the same billion rows in arrival order from many producers, so timestamps are interleaved. Every chunk of a million rows contains samples from across the whole year. Every chunk's `[min, max]` therefore spans nearly the entire domain, every chunk overlaps every filter, and the reader falls back to reading all thousand chunks. The file format, the statistics and the query are unchanged; only the write order changed, and the entire benefit disappeared. The practical consequences: - **Clustering is a write-path decision.** Sorting or partitioning by the column that filters most queries is what buys skipping; nothing on the read path can recover it afterwards. - **Only one order is free.** Rows have one physical order, so one column gets strong skipping. Others benefit only to the degree they correlate with it — a monotonically issued identifier correlates well with time, an unrelated hashed key correlates not at all. - **Chunk size is a granularity knob.** Smaller chunks give tighter ranges and finer skipping but more footer metadata and shorter runs for the encoders; larger chunks give the opposite. - **Statistics are not an index.** They summarise a block; they cannot locate one row. A filter selecting a single scattered row still reads the chunks that might hold it. ## Diagnosing it in practice When someone reports that a filtered query still reads the whole dataset, the checklist runs in this order: does the filter reference the column the data is clustered by; is the predicate one the reader can compare against a range at all, or does it wrap the column in a transformation that hides its value; are the statistics actually present for that column; and is the chunk size so large that even good clustering leaves each range spanning too much. The most common finding is the first — the data is clustered by one column and the query filters on another. The answer an interviewer is listening for is not "the footer has min and max". It is the conditional: skipping is a property of **data layout plus query shape**, and the format only provides the mechanism that lets the two meet.
- When do these statistics stop helping altogether?When the filtered column is uncorrelated with the order rows were written in. Every chunk then holds values from across the whole domain, so every recorded range overlaps every filter and no chunk can be ruled out. The same happens when the predicate wraps the column in a transformation, because the reader can no longer compare it against a stored range.
- Why can only one column get strong skipping cheaply?Rows have a single physical order. Clustering by one column makes that column's per-chunk ranges narrow; every other column's ranges are narrow only insofar as it correlates with the clustering key. Getting strong skipping on a second, uncorrelated column means a second copy in a different order, with all the duplication that implies.
- What does a chunk whose range is entirely inside the filter save?The per-row comparison. Every row in it qualifies by construction, so the reader can take the whole block without testing values. The block is still fetched if the query projects that column, so the saving is processor time rather than bytes — smaller than a skip, but free to take.
saying these in an interview costs you the question
- Believes recorded statistics let a reader skip individual rows
- Thinks skipping works regardless of the order rows were written in
- Assumes an overlapping chunk needs no per-row check after reading
- Treats a missing statistic as permission to skip the chunk
- Confuses per-chunk statistics with a secondary index over values