skip to content

Columnar Storage & Encoding

How an analytical engine physically arranges data on disk so a query touches only the columns and blocks it actually needs. This is the real reason a scan over a billion rows finishes in seconds, and interviewers expect the mechanism, not just the word 'columnar'.

on this pageshow

questions

29

In a columnar table, why does selecting 3 columns of 200 read far less data than SELECT *?

level: juniorimportance: must knowfreq 80%

answer

  1. one column's values live together on disk
  2. metadata records where each chunk starts
  3. the scan asks for byte ranges, not rows
  4. columns you never name are never opened

basics

~20 s

Each column's values sit contiguously in their own chunks, and the file metadata records where every chunk starts. The scan reads only the byte ranges of the columns the query names; the other columns are never opened or decompressed.

solid answer

~50 s

A columnar table is cut into horizontal slices (row groups or stripes); inside each slice, every column's values for those rows form one contiguous **column chunk**, and the metadata records each chunk's offset and length. The scan operator applies **projection pushdown**: it resolves the columns the query actually references, looks up their offsets, and issues reads only for those byte ranges — often as ranged GETs against object storage. The remaining 197 columns are never fetched, decompressed or decoded, so bytes read track the *width of the projection*, not the width of the table. `SELECT *` throws that away: it reads every chunk of every slice and decodes values the query then discards. On an engine billed by bytes scanned that is a direct multiplier on the bill; on a fixed cluster it is wasted I/O, memory and CPU. Note that filter, join and group-by columns are read too, even when they are not in the SELECT list.

code

sql · 7 lines
sql
-- reads every column chunk of every row group it visits
SELECT * FROM events WHERE event_date = DATE '2026-01-01';

-- reads only the three named chunks plus the filter column
SELECT user_id, event_type, amount
FROM events
WHERE event_date = DATE '2026-01-01';

go deeper

for a junior

Be ready to say that values of one column are stored together, so a query reads only the columns it names. Know that SELECT * defeats this and that bytes read scale with the projection.

for a middle

Explain the mechanics: row groups, per-column chunks, the offset metadata that lets the reader jump straight to a chunk, and projection pushdown as the optimizer step that hands the scan its column set.

for a senior

Show you can quantify it — quote the engine's bytes-scanned figure for both query shapes, trace a runaway cost back to a BI connector issuing SELECT *, and reason in compressed bytes per column rather than column counts.

for a principal

Own the governance angle: exposing wide tables through projecting views, treating SELECT * in scheduled jobs as a defect, and understanding that on bytes-scanned pricing the projection discipline is a budget control, not a style rule.

## What "columnar" means physically A columnar table does not store rows end to end. It is first cut into horizontal slices — commonly called row groups, stripes or blocks, typically tens of thousands to a few million rows each. Inside one slice, all the values of a single column are written next to each other as one contiguous byte range: a **column chunk**. A slice of a 200-column table therefore contains 200 chunks, one per column, laid out one after another. Alongside the data, the storage layer keeps metadata that lists, for every slice and every column, at least the chunk's byte offset, its compressed and uncompressed length, and its value count. That offset table is the thing that makes selective reading possible at all: without it the reader would have to walk the bytes sequentially to find where a column begins. ## Projection pushdown When a query is planned, the engine collects the set of columns the query genuinely references — the SELECT list plus anything used in WHERE, JOIN conditions, GROUP BY, ORDER BY, and window definitions — and pushes that set down into the scan operator. This is **projection pushdown**. The scan then consults the metadata, computes the byte ranges for exactly those columns in each surviving slice, and issues reads for those ranges only. On object storage this becomes a set of ranged GET requests; on local storage, a set of seeks. Everything else is not "read and skipped". It is never transferred, never decompressed, never decoded into values. That is the whole economic argument for the layout: an analytical query typically touches a handful of columns out of a very wide table, and the layout makes the cost proportional to what it touches. ## Reason in bytes, not in column counts "3 of 200 columns, so 1.5% of the data" is the right instinct but the wrong arithmetic. Columns are not the same size. A single JSON or free-text column can be larger on disk than a hundred integer columns combined, and encoding widens the gap further: a low-cardinality dictionary-encoded column can shrink to a few bits per row while a high-entropy string column barely compresses. So the useful mental model is: bytes read ≈ the sum of the compressed sizes of the chunks you named. Dropping one wide column from a projection can save more than dropping fifty narrow ones. ## What `SELECT *` actually costs here In a row store, `SELECT *` mostly costs extra network and result-set width — the engine had to fetch the whole row anyway. In a columnar engine it is qualitatively different: it converts a query that touches a thin sliver of the table into one that reads and decodes every chunk of every slice it visits. Consequences in production: - On a warehouse priced by bytes scanned, the invoice scales with the projection. The same logical query can differ by one or two orders of magnitude. - On a fixed-size cluster it burns I/O bandwidth, decompression CPU and memory that other queries needed. - It makes the table hostile to evolution: adding one wide column silently raises the cost of every `SELECT *` consumer, even those that ignore the new column. The usual production source is not hand-written SQL but tooling: a BI connector that pulls all columns and discards them client-side, or an ELT step written as `INSERT INTO ... SELECT * FROM ...`. ## Where the saving does not apply - **Predicate and join columns count.** `SELECT id FROM t WHERE country = 'FR'` reads the `country` chunks as well as `id`. - **`LIMIT` does not narrow the projection.** It caps rows returned; the set of columns opened is unchanged, and with a filter the engine may still scan a lot before it can stop. - **Single-row lookups are not cheap.** Fetching one row by key still means touching one chunk per projected column, so a point lookup pays roughly per column, not per row. Columnar layout rewards narrow reads over many rows, not wide reads of one row. - **Nested data varies.** Whether an engine can read one field of a struct without its siblings depends on how the nested type is shredded into leaf columns. - **Very narrow tables.** With five columns, per-chunk and metadata overhead makes the ratio far less dramatic. ## How to demonstrate it in an interview Run both forms and quote the engine's own reported bytes-read or bytes-scanned figure for each; that number is the direct evidence, and every analytical engine surfaces it. The standard remediation is equally concrete: name columns explicitly, expose wide tables through views that project only the used columns, and treat a `SELECT *` in a scheduled job or dashboard as a defect rather than a style preference.

  • Does a LIMIT 100 on a SELECT * make it cheap again?
    No. LIMIT caps the rows returned, not the columns opened — the scan still resolves every column's chunk in whatever slices it reads. With no filter an engine may stop after the first row group, but with a selective predicate it can read a great deal before it accumulates 100 rows, and it decodes all columns while doing so.
  • If a query projects two columns but filters on a third, what does the scan read?
    All three. Projection pushdown uses the columns the query references anywhere — SELECT list, WHERE, JOIN keys, GROUP BY, ORDER BY — not just the output list. The filter column is read first so the engine knows which rows survive; whether it then reads the projected columns for all rows or only surviving positions depends on whether the engine materializes early or late.
  • Why is a single-row lookup by primary key not a strength of this layout?
    Because cost is paid per column, not per row. Retrieving one row with twenty projected columns means locating and decoding twenty separate chunks, each of which may need a whole compression block decompressed to expose one value. Narrow scans over many rows are the layout's sweet spot; wide reads of one row are the case a row store handles better.

A columnar table is a filing cabinet with one drawer per field rather than one folder per record. Naming three fields opens three drawers; SELECT * means pulling open all two hundred and carrying them to the desk.

saying these in an interview costs you the question

  • Says the whole file is read and then filtered in memory
  • Claims columnar storage makes every query faster, including point lookups
  • Forgets that WHERE and JOIN columns are read even when not projected
  • Thinks adding LIMIT reduces the columns a scan opens
  • Estimates savings by counting columns while ignoring column width

context

open as a page

In a columnar table, how does a per-column encoding differ from a block codec like ZSTD?

level: juniorimportance: must knowfreq 50%

basics

~20 s

An encoding rewrites one column's values in a type-aware way — dictionary codes, run lengths, deltas — and can often be queried directly. A codec such as ZSTD then compresses those encoded bytes opaquely, and the block must be decompressed before anything can read it.

open as a page

Why is inserting rows one at a time into a columnar analytical table so much worse than batching?

level: juniorimportance: must knowfreq 70%

basics

~20 s

Columnar engines write immutable files, and the smallest unit a write can produce is a whole file. One row per insert means millions of tiny files, terrible compression, metadata bigger than data, and a background merge queue that can never catch up.

open as a page

In a columnar engine, what is a zone map and how does it let a query skip data?

level: juniorimportance: must knowfreq 72%

basics

~20 s

A zone map is small per-block metadata holding each column's minimum and maximum value in that block. Before reading a block, the engine tests the query's filter against those bounds and skips every block that cannot possibly contain a matching row.

open as a page

Why does loading a columnar table in sorted order shrink it far more than loading it randomly?

level: middleimportance: must knowfreq 65%

basics

~20 s

Sorting puts equal values next to each other, so run-length encoding sees long runs and delta encoding sees small differences instead of noise. The same data can compress several times better sorted, and correlated columns downstream of the sort key gain too.

open as a page

In a columnar table whose data files are immutable, how does an UPDATE or DELETE actually work?

level: middleimportance: must knowfreq 75%

basics

~20 s

The engine never edits a file. It finds the files containing affected rows, writes new files holding the surviving and modified rows, and atomically swaps the file set in the table metadata. Old files stay until retention expires, then are dropped.

open as a page

Why does filtering a time-sorted columnar table by user_id prune almost no blocks?

level: middleimportance: must knowfreq 66%

basics

~20 s

Because 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.

open as a page

In a columnar engine, why does a contiguous column vector let the CPU filter values with SIMD instructions?

level: middleimportance: must knowfreq 55%

basics

~20 s

Because one column's values sit adjacent in memory at identical fixed width and type, a single SIMD instruction can compare eight or sixteen of them per cycle, with no per-value type dispatch, pointer chasing or branching.

open as a page

What is late materialization in a columnar engine, and when does it beat fetching projected columns up front?

level: seniorimportance: must knowfreq 58%

basics

~20 s

Late 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.

open as a page

In a columnar engine, why can a filter written as a row-wise scalar UDF run an order of magnitude slower than the same logic in built-in operators?

level: seniorimportance: must knowfreq 50%

basics

~20 s

A row-wise UDF is a black box invoked once per value, so it collapses batch-at-a-time execution back to row-at-a-time: no SIMD, no branch-free evaluation, often a cross-language boundary per call, and the optimizer can neither push it to storage nor reorder it confidently.

open as a page

Why do columnar formats split each column into row groups rather than one contiguous run per column?

level: middleimportance: should knowfreq 55%

basics

~20 s

Row groups bound the memory a writer and reader need, give the engine an independently readable unit it can hand to a parallel worker, and give each column chunk its own metadata so whole chunks can be skipped.

open as a page

In a columnar table, how does an engine rebuild a full row when each column is stored separately?

level: middleimportance: should knowfreq 50%

basics

~20 s

By ordinal position. Within one row group every column chunk holds the same rows in the same order, so the k-th value of each chunk belongs to the k-th row. No row identifier is stored beside the values; the shared ordering is the join.

open as a page

Why do columnar engines target data files of hundreds of megabytes rather than a few megabytes?

level: middleimportance: should knowfreq 55%

basics

~20 s

Large files amortise per-file overhead — metadata, catalog entries, storage round trips, task scheduling — and give encodings enough rows to compress well. Files that are too large hurt too: less parallelism, coarser skipping, and every mutation rewrites more bytes.

open as a page

A columnar table is sorted by (event_date, country) — why does filtering only on country prune poorly?

level: middleimportance: should knowfreq 56%

basics

~20 s

Composite sorting is lexicographic: country is ordered only within rows sharing the same event_date. Blocks generally span many dates, so each one contains countries from across the whole list and its country min/max covers nearly everything.

open as a page

In a vectorized engine, how does the batch size affect CPU cache behaviour, and why not use very large batches?

level: middleimportance: should knowfreq 38%

basics

~20 s

Batches are sized so that the vectors an operator touches stay resident in L1 or L2 cache between operators. Too small and per-call overhead dominates; too large and intermediates spill to slower memory, so each operator re-reads its input from RAM.

open as a page

Why can a columnar engine group by a dictionary-encoded string column faster by hashing the integer codes?

level: middleimportance: should knowfreq 40%

basics

~20 s

Grouping on small fixed-width dictionary codes replaces string hashing and byte-by-byte comparison with integer operations on values already in cache. When the code range is small the engine can even index an array directly instead of hashing, and decode group labels once at the end.

open as a page

What happens when a columnar engine dictionary-encodes a column with millions of distinct values per block?

level: seniorimportance: should knowfreq 45%

basics

~20 s

The dictionary grows toward the size of the column itself and each code widens to roughly log2 of the cardinality, so the saving collapses. Most engines detect this and fall back to plain storage, leaving only the general codec to compress the column.

open as a page

When would you choose ZSTD over LZ4 for a columnar table's blocks, and what does it cost?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Choose ZSTD when scans are bound by bytes moved — remote object storage, network, or a storage bill — and LZ4 when scans are already CPU-bound on locally cached data. ZSTD buys a better ratio and charges compression CPU, mostly at write time.

open as a page

How can a columnar engine filter and group a dictionary-encoded column without decoding it?

level: seniorimportance: should knowfreq 42%

basics

~20 s

It evaluates the predicate once per dictionary entry rather than once per row, turning the filter into a set of integer codes, then scans and groups on those codes and decodes only the surviving result values. Run-length encoding lets it aggregate whole runs at once.

open as a page

A columnar table on object storage got tiny appends every minute for a year — why did queries slow down?

level: seniorimportance: should knowfreq 65%

basics

~20 s

The table accumulated hundreds of thousands of tiny files. Planning must enumerate them all, each scan task pays a storage round trip for a few kilobytes of data, and compression is poor — so the query becomes metadata and latency bound rather than throughput bound.

open as a page

What is clustering depth in a columnar table, and how does it show that pruning has degraded?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Clustering depth is the average number of blocks whose key ranges overlap at a given point of the key domain. Depth near 1 means disjoint ranges and near-perfect skipping; depth that climbs over time means new or rewritten blocks now span wide ranges and every query reads more.

open as a page

How would you set ingest batching and compaction policy for a continuously loaded, dashboard-queried columnar table?

level: principalimportance: should knowfreq 40%

basics

~20 s

Start from the freshness the dashboards actually need, size ingest batches to that interval, and let compaction absorb whatever fragmentation remains. Then bound write amplification by choosing how many merge passes data goes through and how large the final files get.

open as a page

Your columnar fact table allows one physical order but three teams filter on different columns — how do you decide?

level: principalimportance: should knowfreq 40%

basics

~20 s

Rank the access patterns by bytes scanned times query frequency, give the physical order to the biggest, then fund the others separately: a partition axis, block-level skipping structures, a second ordered copy, or a pre-aggregate. Decide from measured traffic, not from opinion.

open as a page

When choosing an analytical engine, how much weight should per-core execution efficiency get versus simply adding more compute?

level: principalimportance: should knowfreq 30%

basics

~20 s

Weight it by how much of your workload is actually CPU-bound after pruning. Efficiency buys lower cost per query, better tail latency and more concurrency per node, but scale-out fixes throughput and cannot fix single-query latency or a bad data layout.

open as a page

How do delta, frame-of-reference and bit-packing encodings compress an integer column?

level: middleimportance: nice to knowfreq 38%

basics

~20 s

Delta stores each value as the difference from its predecessor, frame-of-reference stores one per-block minimum plus offsets from it, and bit-packing then writes those small numbers in just enough bits. Composed, a block of rising 64-bit ids can drop to a few bits per row.

open as a page

In an immutable-file columnar table, what is a delete vector and when does it beat rewriting the file?

level: seniorimportance: nice to knowfreq 35%

basics

~20 s

A delete vector is a small side file marking which row positions in a data file are logically deleted. Readers filter them out on the fly. It turns an expensive file rewrite into a tiny write, at the cost of extra work on every read until compaction.

open as a page

How does a per-block bloom filter help a columnar engine where min/max zone maps cannot?

level: seniorimportance: nice to knowfreq 34%

basics

~20 s

Min/max bounds only prove a value falls outside a block's range. A bloom filter answers possible membership, so it can eliminate blocks whose range covers the value but that do not actually contain it — rescuing equality lookups on columns scattered across the table.

open as a page

In a vectorized engine, what is a selection vector and why keep one instead of compacting the batch after each filter?

level: seniorimportance: nice to knowfreq 30%

basics

~20 s

A selection vector is a small list of the positions in a batch that survived a filter. Carrying it lets later operators skip non-matching rows without physically copying every projected column, which would cost more than the filter itself.

open as a page

In a columnar warehouse, when does a 400-column wide table stop being cheap for its readers?

level: principalimportance: nice to knowfreq 38%

basics

~20 s

Width is nearly free only while every consumer projects narrowly. It stops being free once tools issue star-selects, once per-column metadata and short chunks add up, and on the write side, where every column is encoded on every load.

open as a page