The same query: 30 minutes vs 200 ms — column stores read only the columns touched

The same SELECT against the same row count runs ~9000x faster on a columnar engine because the row store reads the whole table while the column store reads only the columns the query touches.

Scene 01

The same query: 30 minutes vs 200 ms

  1. Watch
  2. Try it
  3. Predict
  4. Capture
QUERYSELECT country, avg(latency) FROM events WHERE ts > now()-1h GROUP BY countryNARROW · 3 / 30Postgres (row store)~9000x gapreads ALL 30 columns per rowBYTES READ360 GBWALL CLOCK30:00ClickHouse (column store)reads 3 of 30 column filesBYTES READ28 GBWALL CLOCK0:00.200Narrow query, wide table — the row store hauls all 30 columns; the column store reads only 3.
row store — every column of every row sits together on disk
column store — one file per column, read only what the query asks for
this 9000x gap is what 'OLAP' is named after — Online Analytical Processing
What to watch for

Same query, two engines. The row store pulls 360 GB off disk and spins to 30 minutes. The column store pulls 28 GB and finishes at 200 ms. Same hardware, same row count — only the on-disk shape differs.

Implementation

Highlighted lines are the ones running in the diagram right now.

RowStore.scan(query)
every page holds all N columns — pays the full row width
def scan(query): # Postgres-shaped
bytes_read = 0
for page in heap.pages(): # 8 KB pages, row-major
for row in page.rows(): # all N cols interleaved
bytes_read += row.width # = sum(width(c) for c in N)
tuple = decode(row) # 30 fields materialized
project = [tuple[c] for c in query.columns]
emit(project) # K-of-N kept, N-K discarded
return bytes_read # ~= table_size_bytes
ColumnStore.scan(query)
open one file per queried column — skip the rest entirely
def scan(query): # ClickHouse-shaped
bytes_read = 0
files = [open(f'{c}.bin') for c in query.columns] # K files
# the other (N - K) column files are never opened
for granule in zipGranules(files): # vector of K cols
bytes_read += sum(len(b) for b in granule)
emit(granule) # already projected
return bytes_read # ~= (K / N) * table_size
compare(query, table)
the K/N ratio that the slider is sweeping
def compare(query, table):
N = table.totalColumns # 30 or 100
K = len(query.columns) # 3, 30, or 1
row_bytes = table.size_bytes # always full table
col_bytes = table.size_bytes * (K / N)
ratio = row_bytes / col_bytes # = N / K
return (row_bytes, col_bytes, ratio)

Where this sits in Build a columnar OLAP store (ClickHouse / Druid style)

Scene 01 of 13, in the Why columnar? act — The 30-min Postgres query vs the 200ms ClickHouse query — what's on disk?. Same SELECT, same rows. ClickHouse runs 9000× faster than Postgres because the column store reads only the columns the query touches.

Up next. If both engines hold the same 1.2 billion rows, the only thing that can explain a 9000x gap is what each engine actually pulls off disk — so we need to look at the physical layout under the rows.

All 13 scenes in Build a columnar OLAP store (ClickHouse / Druid style) · Every curriculum

Built with Arqly
Every scene in Build a columnar OLAP store (ClickHouse / Druid style) builds on the one before it.All 13 Build a columnar OLAP store (ClickHouse / Druid style) scenes