OLAP¶
Version and source policy
Engine syntax and feature support are version-sensitive. Check Versions & Primary Sources and reproduce claims on the pinned lab where available.
2:47 PM. A PM wants a live dashboard — p95 latency by endpoint, last hour, one customer at a time — and someone wires it straight to the Postgres replica that already serves the app. Ten minutes later the replica is pegged at 100% CPU on a single GROUP BY, and unrelated app queries start timing out.
Predict before you read on: does this get fixed by (A) a covering index, (B) another read replica, (C) Redis in front of the query, or (D) copying the data into a purpose-built OLAP engine?
Hundreds of millions of product and observability events land every day, and the access pattern behind that dashboard — scan a lot of rows, touch a few columns, aggregate, repeat — is exactly what OLAP engines exist for, not what a row store like Postgres was built to survive.
This is SaaSCo at Stage 7
SaaSCo: The Evolving Company adds ClickHouse once a customer-facing dashboard needs sub-second answers — a latency class Iceberg and Trino were never built to hit, and the lakehouse from Stage 6 keeps serving the cold, historical path.
Workload¶
Two shapes show up everywhere in this academy.
SaaS analytics / observability dashboards
- Ingest: 10⁸–10⁹ events/day.
- Query:
GROUP BY service, endpointover minutes to days, sometimes filtered to onecustomer_id. - Latency: hundreds of milliseconds, not tens of seconds.
- Concurrency: tens (internal) to thousands (customer-facing).
Not this module: federated Iceberg ⨝ Postgres for an analyst. That is Trino. You may load Iceberg into ClickHouse. You do not ask ClickHouse to be a lake catalog.
What this module covers¶
| Topic | What you will be able to do |
|---|---|
| Columnar storage | Account for IO, compression, SIMD; know when rows win |
| ClickHouse | Design ORDER BY, parts, granules, marks; ingest in batches; not pretend mutations are OLTP |
| Pinot | Segments, inverted/star-tree indexes, realtime vs offline, when high-QPS dimensional queries need Pinot |
ClickHouse is the flagship: most teams meet OLAP here, and ORDER BY is a physical design decision you will live with for years. Pinot is the other attractor: user-facing, concurrent, fresh.
Why a row store falls over¶
Postgres stores a row as a tuple. To compute avg(latency_ms) GROUP BY endpoint it must visit tuples. Each tuple carries trace_id, user_agent, payload_hash — columns the query never named.
flowchart LR
subgraph row["Row store page"]
R1["ts · user · ep · lat · status · ..."]
R2["ts · user · ep · lat · status · ..."]
end
subgraph col["Column files"]
EP["endpoint.bin"]
LAT["latency_ms.bin"]
end
Q["GROUP BY endpoint, avg(latency)"]
Q -.->|reads almost everything| row
Q -->|reads two columns| col At 10⁹ rows × 1 KB/row you are reading ~1 TB to answer a two-column aggregate. Columnar layout plus compression plus vectorised loops is why the same query is 100–1000× cheaper in ClickHouse than in OLTP Postgres. The physics lesson is columnar storage.
That speed is bought with constraints: append-heavy writes, batch inserts, expensive mutations, poor point-lookup behaviour if you use the engine as a KV store.
Two engines, two concurrency stories¶
| ClickHouse | Pinot | |
|---|---|---|
| Mental model | Columnar database you own | Realtime OLAP serving layer |
| Index story | Sparse primary index from ORDER BY + skip indexes | Inverted, sorted, range, star-tree, json, text |
| Freshness | Seconds–minutes (batch parts, Kafka MV) | Seconds (consuming segments) |
| SQL | Very wide | Narrower; joins are not the point |
| Concurrency | Good; not “10k QPS dimensional” without care | Built for high QPS on known dimensions |
| Ops | Server + Keeper | Controller, broker, server, minion |
ClickHouse wins when the query is a heavy scan or a rich SQL shape: “p99 latency for this service, excluding health checks, with a HAVING on volume, joining a small dimension.” Internal observability platforms usually land here.
Pinot wins when thousands of tenants each hit a dashboard of the same dimensional queries with fresh Kafka data: “my company’s events, last 15 minutes, broken down by the dimensions we indexed.” Star-tree is pre-aggregation as a first-class index.
If you need both (internal ad-hoc and customer-facing tiles), that is two serving paths, not one compromise cluster. Details: ClickHouse vs Pinot, Pinot.
Physical design is the product¶
In OLTP you add a B-tree when a query is slow. In ClickHouse the sort key is the table. Granules, marks, and the sparse primary index all derive from ORDER BY. Changing it means a new table and a copy.
Pinot is the same idea with more knobs: which column is sorted in the segment, which inverted indexes exist, which star-tree dimensions you pre-aggregate. Those are not “tuning.” They are the schema.
A useful test: write the three queries that must be fast. If they do not share a prefix of the sort key / star-tree dimensions, you need two tables (or a projection / MV), not a wider ORDER BY.
Where OLAP sits¶
flowchart LR
K[Kafka] --> F[Flink / Spark]
F --> I[Iceberg]
F --> CH[ClickHouse]
K --> P[Pinot realtime]
I --> T[Trino]
I --> CH
CH --> G[Grafana / product UI]
P --> U[User-facing dashboards]
T --> A[Analysts] - Lake (Iceberg): source of truth, cheap, slow-ish SQL via Trino.
- OLAP store: copy or stream shaped for dashboards.
- Do not make ClickHouse the only copy of 7 years of events unless you have thought about object storage, TTL, and restore.
Ingest is a batching problem¶
Dashboard engines want parts and segments, not one HTTP INSERT per click. The same Kafka topic can feed:
| Path | Mechanism | Freshness |
|---|---|---|
| ClickHouse | Kafka engine + MV, or Flink/Vector batches of 10k–100k | seconds–minutes |
| Pinot | Realtime consumer → CONSUMING segment | seconds |
| Lake | Flink/Spark → Iceberg commits | minutes–hours |
If the application writes row-by-row to ClickHouse because “we have a JDBC driver,” you will meet too many parts before you meet a slow GROUP BY. That failure is in ClickHouse; the module-level rule is: OLAP ingest is bulk.
CDC (Postgres → Debezium → engine) is the other trap. Event facts append. Customer dimensions upsert (ReplacingMergeTree / Pinot upsert). Treating every OLTP UPDATE as ALTER TABLE UPDATE rewrites columnar parts for a living.
What this module is not¶
| Problem | Go here instead |
|---|---|
| Federated Iceberg ⨝ Postgres | Trino |
| Prom scrape + paging | TSDBs |
user_id as a Prom label | Cardinality — then store the event in CH/Pinot |
Point lookup GET /trace/:id | Search / OLTP / trace backend |
| 12-hour shuffle ETL | Spark |
OLAP is the serving copy of facts you already have (or the primary store if you accepted that operational bargain). It is not a lake catalog and not a queue.
Debugging preview¶
When a dashboard is slow, the first three numbers:
- Rows read vs rows in the time range — sort key / partition prune missed.
- Bytes read vs columns named —
SELECT *or a JSON blob. - QPS × scan — you are on the wrong engine (ClickHouse scan vs Pinot indexed).
ClickHouse: EXPLAIN indexes = 1, system.parts, system.query_log. Pinot: broker scatter vs merge, star-tree hit, consuming lag. Columnar physics: columnar storage.
Scale cliffs¶
| Events / day | Typical pain |
|---|---|
| 10⁸ | Batch inserts, ORDER BY mistakes, too many parts from the application writing row-by-row |
| 10⁹ | Merge pressure, partition explosion (PARTITION BY timestamp instead of month/day), hot customer in the shard key |
| 10¹⁰+ | You shard; you TTL raw data; you pre-aggregate; you stop using FINAL on the hot path |
10× more events is rarely 10× more hardware if the sort key matches the dashboard. It is easily 100× more hardware if every query is a full scan of ORDER BY (timestamp) filtered by customer_id.
How to study this module¶
- Columnar storage — IO for a two-column
GROUP BY, then when columns lose. - ClickHouse — work the observability
ORDER BYexamples with real predicates. Use the ORDER BY explorer. - Pinot — same events, different concurrency and index story.
- Only then read ClickHouse vs Trino and ClickHouse vs Pinot.
Check your understanding¶
Grafana for an observability product: 400 million events/day, 2,000 customers, p95 dashboard 300 ms, peak 80 QPS internal. PM wants the same charts in the customer app at 5,000 QPS, each tenant seeing only their rows, data < 15 s old.
Do you scale the ClickHouse cluster 60×, add Pinot, or put Redis in front of ClickHouse? What physical design must exist either way?
Answer
Do not 60× ClickHouse and hope. Internal 80 QPS of moderately rich SQL is a ClickHouse-shaped load. 5,000 QPS of tenant-scoped dimensional tiles with 15 s freshness is a Pinot-shaped load (or a pre-aggregated ClickHouse path plus aggressive caching, which you will still have to design as a serving schema).
Redis in front of raw event queries fails on the key space (every tenant × every time range × every breakdown) and on freshness.
Either way you must physically cluster by tenant: ClickHouse ORDER BY (customer_id, timestamp) or (customer_id, service, timestamp); Pinot tenant filter with inverted/sorted indexes and likely a star-tree on the dashboard dimensions. Without that, every query scans everyone else's events.