Skip to content

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

0 of 2 checks done
  1. 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.

  2. 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)?
    MB

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

Write at least 60 characters (0 so far). Write it as you would say it in a design review. You will compare it against the points a strong answer makes.

Where you practise it

Further reading

Engineers describing it in systems they run.

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