Joins — denormalize or pay

Column-store joins are bound by RAM (hash join) or network shuffle (sharded), so the canonical OLAP answer is to denormalize dimensions into the fact row at write time — storage is cheap, query latency is the constraint.

Previously

Materialized views collapse one big scan into a tiny one. But the reader's other OLTP instinct — JOIN across normalized tables — hits a different kind of wall on a column store, and the answer feels wrong.

Scene 12

Joins: denormalize or pay

  1. Watch
  2. Try it
  3. Predict
  4. Capture
Star schema · hash joinDenormalized wide tableDictionary lookup · dictGet()events (1B, narrow)users (1M, dim)HASH TABLE (RAM)200 MBbuild from usersprobe per event rowQUERY LATENCY4.0 sRAM: 200 MBright: 1M rowsruntime join: build cost paid once per query, probe paid once per fact rowHash table built on the RIGHT side (users); probed once per event row — RAM-bound by the dim, not by the SQL.
↑ hash table built on the RIGHT side (users) — memory tracks dim size
→ probe per event row — ~4 s wall-clock
What to watch for

The star-schema panel runs the OLTP-instinct query: hash table builds from the 1M-row users dim (200 MB in RAM), then events probes it row by row. Wall-clock pins near 4 s — the hash join is RAM-bound by the right side, not by how cleverly you wrote the SQL.

Implementation

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

hash_join_star_schema()
build hash table on RIGHT (users), probe per LEFT (event) row
def hash_join_star_schema():
# algorithm = 'hash' (ClickHouse default)
ht = {} # hash table in RAM
for u in scan(users): # build on RIGHT
ht[u.id] = u.tier # ~1M rows -> ~200 MB
out = Counter()
for e in scan(events): # probe per LEFT row
tier = ht.get(e.user_id) # one lookup per event
out[(e.country, tier)] += 1
return out # ~4 s; RAM-bound on |users|
denormalized_query()
no hash table; only column reads off a wide fact table
def denormalized_query():
# tier was baked into the row at write time
out = Counter()
for e in scan(events_wide, # only the columns we need
cols=['country', 'tier']):
out[(e.country, e.tier)] += 1
return out # ~0.2 s; one wide scan
# update cost: a user upgrading tier => rewrite every
# event row that user ever produced (parts are immutable).
dictionary_lookup_query()
in-RAM dict pre-loaded once; dictGet per row replaces the JOIN
def dictionary_lookup_query():
# users_dict pre-loaded once into RAM (~30 MB)
# algorithm = 'direct' (no hash-build at query time)
out = Counter()
for e in scan(events, cols=['country', 'user_id']):
tier = dictGet('users_dict', 'tier', e.user_id)
out[(e.country, tier)] += 1
return out # ~0.4 s
# refresh: LIFETIME(MIN 300 MAX 600) reloads from source
# every ~5-10 min — no event-row rewrites on tier change.

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

Scene 12 of 13, in the Patterns act — Materialized views and the OLAP answer to JOIN — denormalize or dictGet.. Column-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.

Up next. We've now seen every lever — layout, encoding, execution, parts, merge, sort key, sparse index, skip index, MVs, denormalization. The last move is to face a real workload and pick the configuration end-to-end.

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