Columnar storage
Storing each column of a table separately, so analytical queries read only the columns they use and compress them well, at the cost of slow single-row lookups and updates.
Storage & state
Learn it
An analytics query like "signups per day from Germany this quarter" touches millions of rows but only two or three columns. A row store keeps each row's columns together, so it reads every column of every matching row and throws most of it away.
A column store lays data out one column at a time, so a query reads only the columns it names.
Work it out
Rows are 1 KB with 50 columns of about 20 bytes. A query uses 3 columns of 10 million rows. Reading only those columns, about how many megabytes are read (before compression)?- Values in one column look alike, so they compress very well, often 10× or more.
- Data is kept in a sort order, which acts as a coarse index: blocks whose min and max rule them out are skipped.
- Data arrives in batches that become immutable parts merged in the background. One row at a time makes too many tiny parts.
- Updating or deleting a row rewrites the parts containing it: fine occasionally, ruinous as a regular workload.
Check
Which workload suits a column store worst?
Quick reference
The same ideas, condensed for revision.
How it goes wrong
- Too many small inserts
- Each insert creates a part; background merging cannot keep up, and the table slows or rejects writes.
- Frequent updates
- Each update or delete rewrites whole parts. Many of them queue behind each other and starve merges.
- Sort order that does not match the queries
- Filters cannot skip blocks, so every query scans the full column.
- Everything in one JSON column
- Every query must read and parse the whole blob, losing most of the benefit of columns.
Instead, consider
- Row store with indexes (Postgres, MySQL)
- Queries fetch or update individual rows, or the data fits comfortably and queries are selective.
- Pre-aggregated rollup tables
- The questions are known in advance, so totals can be computed as data arrives.
- Columnar files on object storage (Parquet) with a query engine
- Data is huge, queried occasionally, and latency of seconds to minutes is fine.
In practice
- DuckDB
- Embedded, in-process column store, good for one machine's worth of data.
- ClickHouse
- A distributed column store built for real-time analytics on event data.
- BigQuery, Snowflake, Redshift
- Managed warehouses that separate columnar storage from compute.
It assumes
- Queries aggregate many rows but read few columns.
- Data is mostly appended in batches and rarely updated.
- The common filters match the table's sort order.
Explain it in your own words
Where you practise it
Further reading
Engineers describing it in systems they run.
- How we're improving performance by combining persons and events
PostHog · PostHog · Post, Nov 2022
Copying person data onto every event to avoid a join at query time, and what that changes about how merged users appear in old data.
- How we turned ClickHouse into our event mansion
PostHog · James Greenhill · Post, Nov 2021
Why event analytics outgrew Postgres, how they chose a column store, and the mistakes they made running it.
- How to speed up ClickHouse queries using materialized columns
PostHog · Karl-Aksel Puulmann · Post, Oct 2021
Finding out where query time goes (parsing JSON), and pulling hot properties into their own columns without rewriting the table.
Related concepts
- Log-structured storage (LSM trees)
Storage engines that turn every write into a sequential append and merge files in the background: very fast writes, at the cost of compaction, tombstones and more expensive reads.
- Partitioning
Splitting data or work by key so each part is handled independently: scaling out, and giving each key a single owner.
- Append-only logs
Recording changes as an ordered, immutable sequence of facts, from which current state, history and replicas can be derived.