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.
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.
2 · The transpose
Rotate the layout 90°. Now the other 38 never leave disk.
3 · The compression
A column is boring, and boring is what compressors eat.
4 · The underrated one
Same file, same statistics, same query — and only one of them is fast.
5 · Keep this card
The whole thing on one index card.
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:
- The data lake is a bucket of Parquet files. Spark, Trino, Athena, and DuckDB all read Parquet natively, which is why “dump it in object storage as Parquet” became a standard landing zone for analytical data. There’s no public market-share figure behind that word; the defensible version is “the format every one of these tools reads first,” not a measured majority.
- In-memory analytics is Arrow. Apache Arrow is a columnar in-memory format — the counterpart to Parquet, not a replacement for it — designed so tools that agree on the layout can hand each other tables without re-encoding them. Pandas 2.x can use PyArrow-backed column types and hand Arrow data around; Polars uses it heavily for interchange, though its internal engine is its own.
- Machine-learning pipelines store their corpora this way. Filtered, deduped training text and its metadata are often kept in columnar files such as Parquet, because “give me every row where
language = enandquality > 0.8, and only these three columns” is exactly the query a data loader makes, and it’s the query the layout is good at.
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:
- Dictionary encoding. A column of country codes with cardinality 200 doesn’t store the strings — it stores integers 0–199 plus a small dictionary.
"United States"once,73a billion times. - Run-length encoding. A column where values repeat (sorted timestamps, sparse flags, low-cardinality enums) is stored as
(value, count)pairs. A million identical values become one pair. - Bit-packing. If a column’s values fit in 9 bits, store them in 9 bits, not 32 or 64.
- Delta encoding. Sorted timestamps become a starting value plus tiny differences, which then compress further.
- General-purpose codecs. Snappy, LZ4, Zstd on top of all of the above.
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.
- Single-row reads are slow.
SELECT * FROM events WHERE id = 42on a columnar store has to seek into every column file. Row stores destroy columnar at this. If your workload is “fetch this user’s profile,” columnar is the wrong tool. - Updates are awkward. In-place update of a single row means rewriting every column’s encoded block. Most columnar systems implement updates via appends + periodic rewrites + tombstones, like an LSM tree — which is part of why “lakehouse” formats like Delta Lake, Iceberg, and Hudi exist on top of Parquet. They also carry transactional metadata, snapshots, and concurrent-writer coordination, which Parquet alone has no notion of.
- The CPU effect is bigger than people expect. Reading a contiguous run of one type lets the CPU vectorize — process 4 or 8 values per instruction with SIMD. Modern columnar engines (DuckDB, ClickHouse, Velox) lean on this hard. Row stores struggle to, because the next field on the cache line is usually a different type.
- Schema evolution is real. Adding a column to a Parquet dataset is cheap (new files have it, old files don’t, the reader handles missing columns). Renaming or retyping is not. Iceberg/Delta exist partly to make this less painful.
- Write amplification is the cost. To get the compression and pruning wins, you batch. To batch, you accept latency between write and visibility, plus background work to compact small files. There is no free lunch — you traded random-write friendliness for read efficiency.
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.
Famous related terms
- Apache Parquet —
Parquet ≈ columnar file format + row groups + per-column stats + nested types— the de facto on-disk format for analytical data. - Apache Arrow —
Arrow = columnar in-memory format + a layout everyone agrees on— the in-RAM counterpart to Parquet; where two tools both speak it, handing a table between them can avoid a re-encode entirely. - Apache ORC —
ORC ≈ Parquet's older cousin from the Hive lineage— same shape, different file format; widely used in Hadoop-era stacks. - Predicate pushdown —
pushdown = move the filter down to the storage layer— what lets min/max pruning skip entire row groups instead of reading them. - Vectorized execution —
vectorized = process columns in batches with SIMD— the CPU-side payoff that makes columnar engines feel fast even on already-cached data. - Lakehouse formats (Iceberg / Delta / Hudi) —
lakehouse table = immutable data files + a transactional metadata layer— bolt-on ACID, schema evolution, and time travel for columnar files (typically Parquet) on object storage.
Going deeper
- Stonebraker et al., C-Store: A Column-oriented DBMS, 2005 — answers “how much of the win comes from the layout and how much from the compression?”, argued at a time when the answer wasn’t yet obvious.
- Andy Pavlo’s CMU 15-721 (advanced database systems) lectures, free online — the best end-to-end answer to “how does an engine actually execute a query over this layout,” with separate lectures on storage models, compression, and vectorization.
- The Apache Parquet format documentation — the rabbit hole for when you need to know exactly which statistics and page indexes are available to prune on, because that determines what your engine can skip.