Heads up: posts on this site are drafted by Claude and fact-checked by Codex. Both can still get things wrong — read with care and verify anything load-bearing before relying on it.
why → how

Why columnar storage won for analytics

Row stores read every column to answer a one-column question. Columnar stores refuse — and that refusal is what makes Parquet, ClickHouse, DuckDB, and every modern data warehouse fast.

Data intermediate Apr 29, 2026 · updated Aug 25, 2026 · 13 min read

On this page

The picture version

Five pictures for a reader who has never thought about how a table sits on disk, following one query: median page-load time for last week.

1 · The problem

Two columns out of forty, and it took forty seconds.

one row on disk: 40 fields, laid out end to end timestamp latency_ms To reach those two the engine reads past the other 38 — on every row. × a billion rows You pay for the data you didn’t ask for, and that bill is the query.
Nothing about a median is hard. The forty seconds are spent on 38 columns nobody asked for, because the row is the unit on disk — not because the arithmetic is slow.

2 · The transpose

Rotate the layout 90°. Now the other 38 never leave disk.

row-oriented tsuserctryms 04-29aliceUS142 04-29bobDE88 04-30carolUS201 the row is the unit — the two you want are scattered columnar tsuserctryms 04-29aliceUS142 04-29bobDE88 04-30carolUS201 the column is the unit — read two aisles, skip thirty-eight Roughly a 20× cut in bytes read, before compression or skipping. But you just made one-row writes expensive. an INSERT now touches forty places on disk instead of one, so formats batch rows into row groups that trade is permanent — it is why this layout lost transactions and won analytics
The same four values, laid out the other way. The inversion is the entire trade: assembling one whole record is now the expensive operation, and answering a two-column question over a billion rows is the cheap one.

3 · The compression

A column is boring, and boring is what compressors eat.

one type, similar values — so formats stack encodings a row store can’t use dictionary a column of 200 country codes stores integers 0–199 plus one small dictionary run-length repeating values become (value, count) pairs — a million identical values, one pair bit-packing if a column’s values fit in 9 bits, store them in 9 bits, not 32 or 64 delta sorted timestamps become a start value plus tiny differences, which compress further then a codec Snappy, LZ4 or Zstd on top of all of the above Still not enough: you asked for last week, and you are decompressing all 52.
In a row, the type changes every few bytes and the entropy jumps; in a column it doesn’t. Compression on a column is qualitatively better, not just more of it — though how much better depends brutally on the data, which is why the usual 5–20×-versus-CSV figure has to be a range.

4 · The underrated one

Same file, same statistics, same query — and only one of them is fast.

the file is split into row groups, each carrying min/max stats per column timestamp ≥ last week written sorted by timestamp wk 1wk 2wk 3wk 4wk 5 wk 6 each group covers a narrow range, so the rest are ruled out without being read or decompressed written in arrival order 1–61–61–61–61–61–6 every group spans the whole range, so every max clears the filter and nothing is skipped The statistics are all still there. They just don’t rule anything out. The only difference is the order the rows were written in. which is why warehouses care so much about clustering, partitioning and sort keys Columnar is only fast if your data is laid out to be skippable.
When a row group’s timestamp.max falls before the filter, the whole group — every column, not just timestamp — is skipped without being read or decompressed. That only happens if the rows inside it are similar, which is a property of how they were written, not of the format.

5 · Keep this card

The whole thing on one index card.

columnar storage = transpose the table on disk + compress each column + skip what the query doesn’t need + a physical sort order you choose the first three happen whether you think about them or not; only the last one is your job
Picture to keep: a warehouse where the shelves run the other way — one long aisle per item type instead of one bin per order, so you walk the “latency” aisle and never enter the other thirty-nine. Where it breaks: filling a single order now means visiting forty aisles. Columnar isn’t a format you adopt; it’s a layout you maintain.

Why it exists

You open the analytics dashboard at work, pick “last week,” and watch a spinner for forty seconds to see a single number: median page load time. Nothing about that number is hard. The site logged a billion clickstream events, each with 40 columns — timestamp, user_id, country, browser, url, referrer, and 35 others — and the answer you want lives in exactly two of them: timestamp and latency_ms. That query is the running example for this post. Two columns out of forty, and it took forty seconds.

Here’s where the time went. A traditional row-oriented database stores each row contiguously: all 40 fields of event 1, then all 40 fields of event 2, and so on. To compute the median latency, the engine has to read past 38 columns it doesn’t care about for every single row, because rows are the unit on disk. You pay for the data you didn’t ask for, and on a billion rows that bill is the query.

Columnar storage is what you get when you stop accepting that bill. Rotate the layout 90°: store all the timestamps together, then all the latencies, then all the user_ids. The query reads exactly the two columns it needs. The other 38 never leave disk.

This isn’t a new idea — the C-Store paper (Stonebraker et al., 2005) and the MonetDB line before it helped establish the modern academic case, and the design appears in proprietary systems earlier than that. What changed is that open formats made columnar storage ordinary in the data-lake ecosystem, while the major warehouses converged on similar layouts anyway: Apache Parquet and Apache ORC on disk, Apache Arrow in memory, and engines like ClickHouse, DuckDB, BigQuery, Snowflake, and Redshift all built around the same layout.

Why it matters now

If you work anywhere near data, columnar is already under your feet:

You rarely choose columnar by typing it. You choose it by picking a tool whose authors already did.

The short answer

columnar storage = transpose the table on disk + compress each column + skip the columns you don't need

Picture to keep: a warehouse where the shelves run the other way. Instead of one bin per order holding forty different items, there’s one long aisle per item type — walk the “latency” aisle end to end and never enter the other thirty-nine. Like a warehouse, except filling a single order is now the expensive operation: you have to visit forty aisles to assemble one. That inversion is the entire trade, and it’s why this layout won analytics and lost transactions.

Three wins stack. Layout: each column is a contiguous run, so a one-column query reads one-column’s worth of bytes. Compression: values in a column are the same type and often similar, so they compress far better than mixed rows do. Skipping: with per-block min/max statistics, the engine can prove an entire chunk of a column is irrelevant and never decompress it.

Row stores aren’t wrong; they’re just answering a different question — give me this whole record — that columnar stores are bad at.

How it works

Build it from the failure, one constraint at a time. Keep the dashboard query in view: median latency_ms where timestamp is in the last week, over a billion 40-column events.

The transpose — and why one fix isn’t enough

A row layout for a tiny table:

row 1: 2026-04-29 | alice | US | 142
row 2: 2026-04-29 | bob   | DE | 88
row 3: 2026-04-30 | carol | US | 201

The same data, columnar:

timestamp:  2026-04-29 | 2026-04-29 | 2026-04-30
user:       alice      | bob        | carol
country:    US         | DE         | US
latency_ms: 142        | 88         | 201

Now the dashboard query reads two columns instead of forty. The bytes you don’t read are the bytes you don’t pay for — in disk I/O, in memory bandwidth, and (on a cloud warehouse) in literal dollars. In the running example that’s roughly a 20x cut in bytes read, before compression or skipping do anything.

Why it breaks: INSERT of one row now touches forty different places on disk instead of one. The transpose that made reads cheap made single-row writes expensive. Fix: stop writing single rows. Columnar formats accumulate rows in memory and write them out in batches — Parquet calls a batch a row group and sizes it in bytes, with the format docs recommending roughly 512 MB to 1 GB; the row count depends on the writer and the schema. Updates become appends plus a periodic rewrite, the same bargain LSM trees make for a different reason. You have traded write latency for read efficiency, permanently, and that trade is the reason columnar lost the transactional workload and won the analytical one.

Why it still isn’t enough: the two columns you do need are still a billion values each. Reading 20x less of a very large thing is still reading a very large thing.

The compression

Fix: exploit the fact that a column is boring. A column is type-homogeneous and usually value-homogeneous — and homogeneity is exactly what compressors eat. That makes it dramatically more compressible than a row, where every few bytes the type changes and the entropy jumps.

Common encodings columnar formats stack:

Parquet files are commonly reported to land somewhere around 5–20x smaller than the equivalent CSV. It has to be a range rather than a figure: the ratio is brutally dependent on the data — cardinality, sort order, and how repetitive the values are all move it by more than the codec does — so any single number would be a lie. The point is that compression on a column is qualitatively better than compression on a row, because the entropy is lower when types and values are clustered.

Why it still isn’t enough: the dashboard asked for last week. Compression shrank the bytes, but you’re still touching every value in both columns and discarding 51 weeks of them. Worse, you have to decompress them to find that out.

The skipping (this is the underrated one)

Fix: write down enough about each batch to prove you can ignore it. A Parquet file is divided into row groups, and each row group carries per-column statistics: min, max, null count, and optionally a Bloom filter and finer-grained page indexes. When you ask WHERE timestamp >= '2026-04-22', the reader checks each row group’s timestamp.max. If it’s less than 2026-04-22, the entire row group — every column, not just timestamp — is skipped without reading or decompressing.

This is predicate pushdown plus min/max pruning, and it’s where a lot of the real-world speed comes from. A query that selectively touches a small slice of the data can physically read a small slice of the data, instead of streaming the whole table and filtering — assuming the layout cooperates. ClickHouse builds on this idea even more aggressively with sparse primary indexes and “skip indexes.” DuckDB, BigQuery, Snowflake — different implementations, same shape.

Why it breaks — and this is the one that bites people: min/max pruning only really fires when each row group’s min/max range is narrow, and that only happens if the rows inside it are similar. A Parquet file whose events were written in arrival order across many servers has every row group spanning the whole week, so every group’s timestamp.max clears the filter and nothing is skipped. The statistics are all still there; they just don’t rule anything out. Fix: control the physical order. Sort or partition on the column you filter by, and each group ends up covering a small slice of it. That’s why warehouses care so much about clustering / partitioning / sort keys — those are the knobs that decide whether pruning fires at all. The corollary is worth stating flatly: columnar is only fast if your data is laid out to be skippable.

Show the seams

A few things the marketing slides skip.

You started with columnar storage = transpose + compress each column + skip the columns you don't need. What did this post add? — + a physical sort order you have to choose. The transpose and the compression happen whether you think about them or not. The skipping — the one that turns forty seconds into one — only fires if the rows were written in an order that makes each row group’s min/max narrow. Columnar isn’t a format you adopt; it’s a layout you maintain.

Check yourself

Before you go — you convert a table to Parquet and the dashboard query gets 15x faster. Then you add a filter on country = 'DE' and it barely improves at all, even though Germany is 3% of traffic. What’s the likely cause?

Answer

Most likely the rows aren’t clustered by country, so every row group contains a little German traffic. With DE inside every group’s min/max range, no group can be excluded, and the engine falls back to reading the whole country column and filtering — the filter runs at execution time rather than at skipping time. The fix is physical, not logical: sort or partition by country when writing, and the same query starts pruning. Notice the query text never changed; only the layout did.

And one more — would you expect columnar to help a query like SELECT * FROM events WHERE event_id = 91827364?

Answer

No — this is the case columnar is worst at, and it’s worth being able to say why in two steps. SELECT * defeats the projection win: you need all forty columns, so there’s nothing to skip. And fetching one row means seeking into forty separate column chunks and decompressing a block from each, versus a row store’s single page read. Columnar’s wins are “few columns, many rows”; this query is “all columns, one row” — the exact inverse. If your workload is full of these, you want a row store, or a row store sitting in front of the columnar one.

Going deeper