Build your own columnar OLAP store (ClickHouse / Druid style)
OLTP picks one row by key; OLAP scans a billion rows of one column and asks for a percentile. Build the analytical engine that makes that fast: columnar layout, dictionary/RLE/delta compression, vectorized execution, late materialization, MPP shuffle. Internalize why Postgres is 1000× slower than ClickHouse on the same query and why the inverse is also true.
- Scenes
- 13 interactive scenes
- Time
- about 91 minutes
- Topic
- Storage Engines & Databases
What you will be able to explain afterwards
- row-store vs column-store (the load-bearing layout choice)
- dictionary + RLE + delta + LZ4 compression by column type
- vectorized execution — process one column-chunk at a time, not one row
- predicate pushdown + zone maps / min-max indexes
- late materialization — defer row reconstruction to after filters
- MergeTree / LSM-on-columns: parts, merges, mutations
- MPP query execution: scan → shuffle → aggregate → final
- sharding by primary key vs by date partition
- approximate queries: HyperLogLog, T-Digest, quantile sketches
- ETL pull (batch) vs streaming ingest with idempotent inserts
Why columnar?
The 30-min Postgres query vs the 200ms ClickHouse query — what's on disk?
- 01The same query: 30 minutes vs 200 ms — column stores read only the columns touchedSame SELECT, same rows. ClickHouse runs 9000× faster than Postgres because the column store reads only the columns the query touches.~7 min
- 02Same table, two on-disk shapes — row-store pages versus per-column .bin filesRow store interleaves a row's columns contiguously; column store stores each column in its own file. Same rows, rotated 90°.~7 min
Speedups
Compression and vectorized execution — where the orders of magnitude live.
- 03Compression — the column store's superpowerAdjacent column values are same-type and often similar, so RLE and dictionary encoding deliver 5–20× shrinkage that fails completely on row pages.~7 min
- 04Vectorized execution — process batches, not tuplesTuple-at-a-time Volcano is interpreter overhead; processing 1024–8192 column values per call lets the CPU emit SIMD inner loops.~7 min
Write side
Bulk inserts → immutable parts → background merge → too-many-parts cliff.
- 05Writes must be bulk, not per-row — per-column insert overhead and async insertsEvery INSERT touches every column file; per-row inserts on a 30-column table pay 30× the per-column overhead. Batches or async_insert.~7 min
- 06A part — one batch, frozen on disk — sparse primary index inside an immutable partEach batched write lands as an immutable directory of column files plus an index — a part. A table is a stack of parts.~7 min
- 07Merge — and the 'too many parts' crashBackground worker fuses small parts into bigger ones; when write rate exceeds merge rate, parts_to_throw_insert fires and inserts get rejected.~7 min
Read side
Sort key, sparse primary index, skip indexes — narrowing what to scan.
- 08ORDER BY — filtering becomes range-scanRows inside a part are sorted by ORDER BY; WHERE on a leftmost-prefix column is a contiguous range, anything else is a full scan.~7 min
- 09Granules and the sparse primary indexRows are grouped into 8192-row granules; one index entry per granule keeps the whole table's index in RAM, binary-searched in microseconds.~7 min
- 10Skip indexes — prove a granule has no match, then skip itPer-granule sketches (minmax / set / bloom) prove 'this granule cannot contain a match' and skip reading it — only useful when values are clustered.~7 min
Patterns
Materialized views and the OLAP answer to JOIN — denormalize or dictGet.
- 11Materialized views are INSERT triggersA column-store MV is not a refresh — it's a trigger over the incoming block. Direct INSERT into the target silently bypasses it and drifts.~7 min
- 12Joins — denormalize or payColumn-store joins are RAM-bound or shuffle-bound; the canonical OLAP answer is to denormalize at write time — storage is cheap, latency is the constraint.~7 min
Design canvas
Pick the workload, configure the engine, watch each failure mode cite its scene.
More in Storage Engines & Databases
Open the box every design diagram labels "DB": pages, logs, LSM trees, wide-column, documents, graphs and columnar scans, built from scratch.
- Build Build a Bitcask-style KV storeThe simplest possible KV store that still works: an append-only log on disk + an in-memory hash index. Build it from first principles and feel which trade-offs every later store inherits.
- Build Build an LSM-tree storage engine (LevelDB / RocksDB style)The simplest possible storage engine that gives you BOTH ordered reads AND more keys than fit in RAM, by accepting a deal: write to RAM at memory speed, log to disk for safety, then merge sorted files in the background forever.
- Build Build a B-tree storage engine (SQLite-style)What actually happens when you run INSERT INTO users(...). One file of fixed-size pages, organized as B-trees, with a write-ahead log that turns commits into appends. Build it from a SQL writer's perspective and feel why every knob exists.
- Build Build a wide-column store (Cassandra / DynamoDB family)One server is not enough — disk fills, throughput maxes, the box dies. Build a multi-node store from first principles: hash sharding, the consistent-hash ring, vnodes, replication, eventual consistency, tunable W+R quorum, hinted handoff, read repair. Every modern Dynamo-style store is a point in this design space.
- Build Build a graph database (Neo4j / Dgraph-style)When the workload is 'friends of friends', a relational join melts. Build a store where edges are first-class — index-free adjacency, traversals that follow pointers instead of joining tables, and a query language (Cypher / GraphQL+) that thinks in patterns. Feel why graph storage shines for traversal-heavy work and stumbles on full-graph aggregates.
Prefer to design it yourself?
The same subject as a staged workspace: draw the architecture, and a simulator traces requests through the boxes you drew.
Open the Build a columnar OLAP store (ClickHouse / Druid style) workspace