Why This Matters#
A real-time OLAP store is not a faster data warehouse. It is a pre-paid answer to a known set of questions. ClickHouse, Apache Druid and Apache Pinot all win the same way: they decide at ingest time how data is sorted, encoded, indexed and partly aggregated so that a GROUP BY over billions of rows comes back in tens of milliseconds while Kafka is still writing to the same table. Every one of them loses the same way: someone asks a question the layout was not built for, the query scans 40 TB instead of 40 MB, and a user-facing dashboard that promised 100 ms p99 times out at 30 seconds.
That is why it keeps showing up in Staff loops. "Design ad click aggregation", "Design the analytics page for creators", "Show every merchant their sales by hour", "Who viewed my profile", "Build a real-time QoE dashboard for video": each one ends at the same box labelled serving store for aggregates. The L5 candidate writes "ClickHouse" in it. The L6 candidate says "the query shapes are fixed: three filters, two group-bys, a time range. I'll sort by (tenant_id, event_time), roll up to the minute at ingest, and accept losing per-event drill-down past 7 days." The L7 candidate asks who is allowed to add a dimension, because every new column on a user-facing analytics product is a permanent cost line.
The gap between levels is not knowing what a columnar format is. It is knowing that the sort key is the index, the rollup is the cost model, and query concurrency, not data volume, decides which engine you pick. A warehouse answers 50 analysts' unpredictable questions in seconds. A real-time OLAP store answers 50,000 users' predictable questions in milliseconds. Confusing the two is the most expensive mistake in this space.
The L5 → L6 → L7 Contrast#
| Behavior | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| First move | "Put the events in ClickHouse" | "Who queries this, how many at once, at what p99? Internal analysts at 20 QPS and merchants at 5,000 QPS are different systems." | "Is this one shared analytics platform or a store per product? That decides cost and headcount for three years." |
| Data model | One wide table with every field | Sort key from the dominant filter, rollup granularity from the product's slowest acceptable drill-down | Writes the dimension-admission policy: who may add a column, at what cost, with what retention |
| Freshness | "Real time" | Seconds of freshness from Kafka for the head; hourly compacted batch segments for the tail; states the lag SLO | Prices freshness: the delta between 5-minute batch and 5-second streaming, per product |
| Correctness | Counts what arrives | Names the dedup and upsert mechanism, and that background merges do not guarantee it | Decides which numbers are billing-grade (warehouse of record) and which are directional (OLAP), and writes that on the dashboard |
| Cost | Node count | Bytes scanned per query × QPS; rollup ratio; storage tiers | $ per 1,000 queries per tenant; chargeback; a kill switch for the dimension that triples the bill |
| Build vs buy | Picks the engine they know | Picks by concurrency, upsert needs, and join needs | Self-host vs managed crossover modelled in engineer-months; exit path for each engine |
Why "First move" separates levels
The three engines overlap 80% on features. What separates them, and what separates a good design from a bad one, is the query population. Fifteen analysts running ad-hoc SQL tolerate 3-second answers, want joins, and change their questions daily: that is a warehouse or ClickHouse with generous hardware. Two million merchants each loading a dashboard with five tiles is 10,000 QPS of nearly identical queries with a 200 ms page budget: that is Pinot or Druid with pre-aggregation and aggressive pruning, or ClickHouse with materialized views and a cache. The Senior answer picks the engine first and discovers the query population in production. The Staff answer asks for it in minute one.
Why "Correctness" separates levels
Every real-time OLAP engine ingests at-least-once from Kafka by default. ClickHouse's ReplacingMergeTree deduplicates only when parts merge, at an unspecified time, and its own documentation says it does not guarantee the absence of duplicates. Pinot upserts are exact but require the stream to be partitioned by primary key and cost heap memory per key. Druid's streaming rollup is best-effort, not perfect. A Staff candidate names which of these they rely on. A Principal candidate goes one step further: the invoice is computed in the warehouse of record, and the OLAP number is labelled estimated until reconciliation closes.
The 60-Second Pitch#
"For user-facing analytics I'd put a real-time OLAP store behind the API: Pinot or Druid if this is thousands of concurrent, predictable queries; ClickHouse if it's fewer, more ad-hoc queries or the team wants plain SQL with joins. Events land in Kafka, keyed by tenant. The store consumes Kafka directly, so freshness is 5–30 seconds. Data is columnar, dictionary-encoded and sorted by tenant then time, so a tenant's query reads megabytes, not terabytes. I'll roll up to one-minute grain at ingest, which cuts rows 20–100× for bounded dimensions, and keep raw events in object storage for 30 days for drill-down and backfill. Distinct counts use HLL sketches. Billing-grade numbers come from the warehouse, not from here. The capacity unit is bytes scanned per query times QPS, not terabytes stored."
That pitch names the query population, the freshness path, the sort key, the rollup ratio, the approximation, the boundary of trust, and the capacity unit. Those are the seven things an interviewer listens for.
The Three Intents#
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| User-facing analytics (merchant dashboards, creator stats, "who viewed my profile") | 1K–50K QPS, p99 < 100–300 ms, every query scoped to one tenant | Pinot/Druid or ClickHouse with materialized views; tenant-first sort key; star-tree or rollups; result cache | One large tenant's scan saturates shared servers; p99 collapses for everyone | Directional; small, labelled lag is acceptable |
| Operational / real-time monitoring (ad pacing, fraud, QoE, canary analysis) | Freshness < 30 s, high-cardinality slicing, tens of QPS from dashboards and alerts | Stream ingest, short hot retention (7–30 days), approximate distinct counts | Ingest lag during a spike hides the incident the dashboard exists to show | Approximate is fine; staleness must be visible |
| Internal ad-hoc exploration (product analytics, log analytics) | 10–100 QPS, unpredictable queries, joins, long retention | ClickHouse or a cloud warehouse; raw events, wide tables, projections | One analyst's unbounded query burns the cluster; nobody owns the schema | Reproducible; exact when it matters |
🎯 Staff Move: "I'll design for user-facing analytics first. That's where concurrency, tenant isolation and a hard p99 budget meet, and it forces every layout decision. If this is really internal exploration with a dozen analysts, I'd stop and recommend the warehouse we already have; a second analytics engine is not free."
The Staff Positions#
| Position | Rationale |
|---|---|
| Design from the query, not the event | List the 5–10 query shapes first. The sort key, rollup grain and indexes fall out of them; a layout built from the event schema scans everything. |
| The sort key is the primary index | All three engines prune by sorted or partitioned columns. tenant_id first turns a 10 TB table into a 50 MB read for one tenant. |
| Roll up at ingest, keep raw elsewhere | Pre-aggregation is the biggest cost lever (often 20–100×), and it is irreversible inside the store, so raw events live in object storage for replay. |
| Bytes scanned × QPS is the capacity unit | Storage is cheap; CPU to scan and decompress is not. Size by the worst common query, not by disk. |
| Exactness is a property you buy, not a default | Ingest is at-least-once. Upserts, dedup and perfect rollup each have a memory or latency price; pay it only where the number is acted on. |
| High-cardinality dimensions are admitted, not added | A user_id dimension destroys rollup and bloats indexes. Each new dimension gets a cardinality estimate and an owner. |
| The OLAP store is a derived view, never the record | It must be rebuildable from Kafka plus object storage. If it cannot be rebuilt in a day, it has become a database you didn't design. |
Architecture & Internals#
Five internals change design decisions: columnar layout and encoding, segments or parts as the unit of everything, the stream ingest path, scatter-gather query execution, and indexes that skip data rather than find rows.
Columnar Storage — Why Scans Are Cheap#
A row store writes (ts, tenant, country, device, campaign, clicks, spend) together. A column store writes each column as its own file, sorted in the same row order. A query that touches 3 of 40 columns reads ~7% of the bytes before compression even helps.
Then compression stacks on top, because a sorted column is full of runs and repeats:
| Encoding | Works On | Typical Effect |
|---|---|---|
| Dictionary | Strings with bounded distinct values (country, device) | 2–4 byte integer IDs instead of strings; enables bitmap indexes |
| Run-length | Sorted, low-cardinality leading columns (tenant_id when sorted first) | A million identical values stored as one run |
| Delta / double-delta | Timestamps, monotonic IDs | A few bits per value |
| Bit-packing | Small integers after dictionary encoding | 7 bits for 100 distinct values instead of 32 |
| General codec (LZ4, ZSTD) | Everything, on top of the above | Another 2–4×; LZ4 for speed, ZSTD for cold data |
A design number worth stating: well-sorted event data commonly lands at 5–15× compression versus raw JSON, and much more for low-cardinality columns. The design consequence is that column order in the sort key changes the bill: sorting by a low-cardinality column first makes every column after it compress better.
🎯 Staff Insight: "Compression is a property of sort order. If I sort by
event_idfirst, every other column looks random and I pay three times the storage and scan cost. Tenant, then a coarse time bucket, then the next most common filter."
Three Engines, One Shape#
All three are shared-nothing clusters of immutable, columnar chunks with a routing tier in front. The differences are in who owns ingest and how much is pre-computed.
| Concept | ClickHouse | Apache Druid | Apache Pinot |
|---|---|---|---|
| Unit of data | Part (immutable, merged in background) inside a partition | Segment (time-chunked, typically 300–700 MB) | Segment (realtime "consuming" or completed, offline) |
| Ingest from Kafka | Kafka table engine + materialized view, or an external consumer doing batched INSERT | Kafka indexing service: supervisor spawns ingest tasks | Servers consume Kafka partitions directly; segments seal and upload to deep store |
| Query routing | Distributed table fans out to shards | Broker → historicals + realtime tasks | Broker → servers, prunes by time, partition, metadata |
| Durability | Replication via ClickHouse Keeper (ZooKeeper-compatible); local disks or object storage | Deep storage (S3/HDFS) is the source of truth; historicals cache | Deep store holds completed segments; servers replicate consuming segments |
| Pre-aggregation | Materialized views into AggregatingMergeTree / SummingMergeTree, projections | Rollup at ingest is a first-class table setting | Star-tree index; ingestion aggregation |
| Mutable data | ReplacingMergeTree (eventual), lightweight deletes, mutations (heavy) | Re-ingest the time interval (segment replacement) | Full and partial upsert on realtime tables; dedup |
| Coordination | Keeper for replication only | ZooKeeper + metadata DB (Postgres/MySQL) | Helix on ZooKeeper |
| Sweet spot | Ad-hoc SQL, joins, log and product analytics, one team running it | Time-series slice-and-dice with rollup, high-concurrency dashboards | Very high QPS user-facing queries, upserts, tenant-scoped reads |
The practical reading: ClickHouse is a fast SQL database you shape with tables; Druid and Pinot are serving systems you shape with ingestion specs and indexes. ClickHouse gives more freedom and fewer moving parts; Druid and Pinot give more built-in machinery for the serve-thousands-of-users case.
The Write Path — Kafka to Queryable in Seconds#
- Consume. Each Kafka partition maps to one consumer. Rows land in an in-memory, row-oriented buffer that is already queryable. Freshness is poll interval + flush, typically 5–30 s.
- Seal. When the buffer hits a row or time threshold, it is converted to a columnar, indexed segment, uploaded to deep storage, and the Kafka offset is committed with it. In ClickHouse, every batched insert becomes a new part directly.
- Serve. Historical servers load sealed segments (memory-mapped). The query layer merges results from consuming and sealed segments.
- Compact. Small segments or parts are merged into larger ones; in Druid and Pinot, compaction can also re-roll minute data to hourly. ClickHouse merges parts continuously in the background.
Three design consequences:
- Small inserts kill ClickHouse. Each insert creates a part. Thousands of tiny inserts per second create more parts than merges can absorb and trigger the "too many parts" error. ClickHouse's own guidance is batches of at least 1,000 rows, ideally 10,000–100,000, or server-side async inserts (ClickHouse bulk insert docs).
- Partition count caps ingest parallelism. One Kafka partition feeds one consumer. A 16-partition topic cannot be ingested by 64 servers. And in Pinot upsert tables, topic partitions cannot be increased after table creation, so choose partition count once.
- Segment size is a query-cost decision. Druid recommends 300–700 MB segments, roughly 5M rows as a starting target (Druid segments docs). Thousands of tiny segments make every query fan out to thousands of tasks.
The Read Path — Scatter, Prune, Gather#
Query latency is max(server latencies) + merge, so the slowest server sets the p99. Two levers matter more than CPU:
- Pruning before fan-out. Brokers skip segments whose time range, partition or min/max metadata cannot match. A query that reaches 12 segments instead of 2,400 is 200× cheaper before any index runs.
- Partial aggregation on the server. Servers return
(country, sum)pairs, not rows.COUNT(DISTINCT user_id)cannot merge this way, which is why sketches (HLL, Theta) exist: they merge.
Indexes That Skip Data#
OLAP indexes rarely find a row. They prove a block cannot match, so it is never decompressed.
| Index | Engine | How It Helps | Cost |
|---|---|---|---|
| Sparse primary index (sort key) | ClickHouse | One mark per granule of 8,192 rows by default; binary search over marks small enough to stay in memory (MergeTree docs) | Only helps filters on the sort-key prefix |
| Sorted column | Pinot, Druid (time) | Range of row IDs for a value: O(log n) | One per table |
| Bitmap / inverted | Druid (default on dimensions), Pinot | Roaring bitmap per value; AND/OR across filters | Grows with cardinality; useless for unique IDs |
| Range | Pinot | Numeric range filters without scanning | Extra storage |
| Bloom filter / skip index | All three | Skip blocks for point-ish lookups (request_id = ...) | False positives; tuned per column |
| Star-tree | Pinot | Pre-aggregated tree over chosen dimensions; bounded work per query | Storage; build time; not usable with upsert |
| Projection / materialized view | ClickHouse | Second sort order or pre-aggregated copy, chosen automatically or by query | Doubles write and storage for that copy |
The Star-Tree Index — Pre-Aggregation You Can Tune#
Pinot's star-tree is the clearest example of trading space for a latency ceiling. It is built over an ordered list of dimensions (dimensionsSplitOrder) and a set of pre-computed aggregations (functionColumnPairs such as SUM__clicks). Each level splits on one dimension; a special star node at each level holds the aggregate across all values of that dimension. Nodes stop splitting when they hold fewer than maxLeafRecords records (default 10,000). Pinot's documentation describes it as a configurable space-time trade-off that gives a hard upper bound on query latency for a use case (Pinot star-tree docs).
A query WHERE device = 'ios' GROUP BY nothing walks country = * then device = ios and reads one pre-aggregated record instead of scanning every iOS row. The two rules that make or break it:
- Every aggregation in the query must be in
functionColumnPairs, or the star-tree is ignored and the query falls back to a scan. - Split order follows filter frequency. Put the dimensions that appear in most
WHEREclauses first.
🎯 Staff Insight: "Star-tree is a rollup that doesn't throw away the raw rows. I get the latency of pre-aggregation for my 10 dashboard queries and still keep row-level data for the rare drill-down. The price is storage and a fixed list of metrics, so new metrics need a rebuild."
Data Modeling — "The Entire Game"#
In OLTP you model entities and let the planner find a path. In real-time OLAP you model the queries, and the physical layout is the plan. A bad sort key cannot be fixed by hardware; it can only be fixed by rewriting the table.
Step 1: Write Down the Query Shapes#
Take a merchant analytics product: 2M merchants, 3B order events/day, a dashboard with five tiles.
| # | Query Shape | Filters | Group By | Freshness | Share of QPS |
|---|---|---|---|---|---|
| Q1 | Sales today vs yesterday | tenant, time (2 days) | hour | < 1 min | 40% |
| Q2 | Top 10 products, last 7 days | tenant, time | product_id | < 5 min | 25% |
| Q3 | Sales by country / channel | tenant, time | country, channel | < 5 min | 20% |
| Q4 | Unique buyers, last 30 days | tenant, time | — (distinct) | < 1 h | 10% |
| Q5 | Order drill-down (raw rows) | tenant, order_id or time window | — | < 1 min | 5% |
Every query is tenant-scoped and time-bounded. That one observation picks the sort key, the partitioning and the routing. Q2 has a high-cardinality group-by (product_id) that rollup will not shrink much. Q4 needs a sketch. Q5 needs raw rows, but only for a short window.
Step 2: Sort Key, Partitioning, Grain#
-- ClickHouse: raw events, short retention, for Q5 drill-down
CREATE TABLE orders_raw
(
tenant_id UInt64,
event_time DateTime64(3),
order_id UInt64,
product_id UInt64,
country LowCardinality(String),
channel LowCardinality(String),
amount_cents Int64,
buyer_id UInt64,
version UInt64 -- for dedup on replay
)
ENGINE = ReplicatedReplacingMergeTree(version)
PARTITION BY toYYYYMM(event_time) -- coarse: never by tenant
ORDER BY (tenant_id, toStartOfHour(event_time), order_id)
TTL toDateTime(event_time) + INTERVAL 14 DAY;
-- Rollup: one row per tenant x minute x country x channel
CREATE TABLE orders_1m
(
tenant_id UInt64,
minute DateTime,
country LowCardinality(String),
channel LowCardinality(String),
orders SimpleAggregateFunction(sum, UInt64),
revenue SimpleAggregateFunction(sum, Int64),
buyers AggregateFunction(uniq, UInt64) -- mergeable sketch state
)
ENGINE = ReplicatedAggregatingMergeTree
PARTITION BY toYYYYMM(minute)
ORDER BY (tenant_id, minute, country, channel)
TTL minute + INTERVAL 13 MONTH;
CREATE MATERIALIZED VIEW orders_1m_mv TO orders_1m AS
SELECT tenant_id, toStartOfMinute(event_time) AS minute, country, channel,
count() AS orders, sum(amount_cents) AS revenue, uniqState(buyer_id) AS buyers
FROM orders_raw
GROUP BY tenant_id, minute, country, channel;
-- Query Q4 reads merged sketch state, not raw buyers
SELECT uniqMerge(buyers) FROM orders_1m
WHERE tenant_id = 42 AND minute >= now() - INTERVAL 30 DAY;
The decisions, in the order a Staff candidate says them:
- Sort key =
(tenant_id, coarse time, ...). The tenant filter becomes a range over marks. ClickHouse's own guidance is to never partition by client identifier, and to put it first inORDER BYinstead; it also says partitioning coarser than a month is rarely needed (MergeTree docs). - Partition by month, not by tenant or day-per-tenant. Partitions are for retention and bulk operations, not for query pruning. 2M tenants as partitions is 2M directory trees and a merge storm.
- Grain = 1 minute for 13 months, raw for 14 days. Product signed off that drill-down to an individual order is only available for two weeks. That sentence is the entire storage budget.
- Materialized view on insert. In ClickHouse an MV is an insert trigger, not a refreshed snapshot: each inserted block is aggregated and written to the target. It never sees data inserted before it was created, so a new MV needs an explicit backfill. It also sees every duplicate insert before
ReplacingMergeTreecollapses it, so a replayed batch is counted twice in the rollup. Dedup must happen upstream of the MV, not in the raw table. - Distinct counts as sketch state.
uniqStatestores a mergeable sketch per row, so 30 days of minute rows merge into one estimate with a small, bounded error (low single-digit percent at most). Exact distinct over 30 days per tenant is a different, much more expensive product.
The Pinot or Druid equivalent is the same plan written as config: tenant_id as the sorted column or partition column, time column with segment granularity, a rollup or star-tree over (country, channel) with SUM__amount_cents, and an HLL or Theta-sketch metric for buyers.
Rollup at Ingest vs Aggregation at Query Time#
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Raw rows, aggregate at query | Full flexibility; any new question answerable; simplest pipeline | Cost scales with rows scanned × QPS; p99 grows with data | Infra budget; users at peak |
| Rollup at ingest (Druid rollup, ClickHouse MV) | 20–100× fewer rows for bounded dimensions; predictable latency | Raw detail gone from the store; dimensions fixed at ingest; streaming rollup is best-effort | Product: no drill-down below the grain; data team for backfills |
| Index-time pre-aggregation (star-tree, projections) | Rollup speed and raw rows kept | Storage 1.2–2× more; fixed aggregation list; rebuild on change | Storage budget; ingest CPU |
| Pre-computed result tables (batch or stream job writes final answers) | Lowest latency; trivially cacheable | Every new question is an engineering ticket | Data engineering team, forever |
Druid's own rollup documentation is blunt about the trade: rollup combines rows with identical timestamp and dimension values, and in exchange you lose the ability to query individual events. Streaming ingestion gives best-effort rollup, where identical keys can remain in several segments until compaction (Druid rollup docs).
The rollup ratio formula:
rollup_ratio = events_ingested / rows_stored
rows_stored ≈ min(events, time_buckets × Π cardinality(dimension_i) actually observed)
Example: 3B events/day, 1-min grain = 1,440 buckets,
active (tenant, country, channel) combos per minute ≈ 400K
rows/day ≈ 1,440 × 400K = 576M → ratio ≈ 5×
Add product_id (avg 6 active per tenant-minute): rows ≈ 3.4B → ratio < 1, rollup useless
That last line is the whole high-cardinality problem in one calculation. Measure the ratio in production (Druid suggests SUM(num_rows) / COUNT(*)), and alert when it falls.
🎯 Staff Move: "Before I add a dimension to the rollup, I multiply its observed cardinality into the row count. If the ratio drops below about 5×, that dimension belongs in a separate table or a star-tree, not in the main rollup."
Upserts and Deduplication#
Real-time OLAP was designed for append-only facts. Production data is not: orders get refunded, Kafka redelivers after a consumer restart, CDC streams emit updates. Each engine handles mutability differently, and each has a price.
| Engine | Mechanism | Guarantee | Price |
|---|---|---|---|
| Pinot upsert (full or partial) | In-memory map from primary key to the latest record location; older versions are masked at query time | Latest version per key, exact, within a partition | Input stream must be partitioned by primary key; heap memory roughly keys × (key_bytes + 24); strictReplicaGroup routing; star-tree unavailable; partitions fixed at creation (Pinot upsert docs) |
| Pinot dedup | Drop records whose primary key was already seen | First write wins, exact | Same partitioning and key-map memory |
ClickHouse ReplacingMergeTree | Keeps the row with the highest version when parts merge | Eventual: dedup happens only during merges, at an unspecified time, and duplicates may remain (ReplacingMergeTree docs) | Query with FINAL or argMax for correct answers, which costs CPU at read |
| ClickHouse insert dedup | Replicated tables drop identical retried insert blocks | Retries of the same block, not semantic duplicates | Retry must resend the same batch |
| Druid | Re-ingest the affected time chunk; new segment version atomically replaces the old | Exact for the replaced interval | Batch job per correction; minutes to hours |
Sizing the Pinot key map is a number to say out loud: 500M live order IDs × (8 + 24) bytes ≈ 16 GB of map spread across servers, before replication. If keys never expire, that grows forever. Set a metadata TTL, or keep upserts only on the hot window and compact the rest to an offline table.
🎯 Staff Insight: "Upsert in an OLAP store is a correctness feature with a memory bill. I'll use Pinot upsert for order status because the dashboard must show refunds within a minute, and I'll keep the key window to 30 days. For everything else I'll make the pipeline idempotent upstream and accept at-least-once with a dedup key."
High-Cardinality Dimensions#
A dimension is high-cardinality when its distinct-value count is close to the row count: user_id, session_id, url, ip, trace_id. Each one does the same three things:
- Destroys rollup. Every row becomes unique; the ratio goes to 1.
- Bloats dictionaries and bitmaps. A dictionary with 500M entries per segment and a bitmap per value cost more than the column they index.
- Explodes group-by state.
GROUP BY urlover a day can produce 100M groups per server, all held in memory before the merge.
What to do instead, ranked by preference:
| Need | Pattern |
|---|---|
| Count distinct users | HLL / Theta sketch metric, mergeable across segments (1–2% error) |
| Top-N URLs | Approximate top-K, or a separate rollup table keyed by URL with daily grain |
| Look up one user's events | Raw table sorted by (tenant_id, user_id, time) with short retention, or the OLTP store |
| Filter by a specific ID | Bloom filter / skip index, not bitmap or dictionary-backed inverted index |
| Group by a long-tail field | Bucket at ingest: top 1,000 values kept, the rest mapped to other |
🎯 Staff Move: "Cardinality is a product question disguised as a schema question. If product needs per-URL analytics, I'll ask how many URLs per tenant and for how long, then give it its own table with its own retention, instead of letting it ride inside the main rollup and multiply its cost."
The Tunable Tradeoff — Freshness × Flexibility × Cost per Query#
You get to pick where each workload sits on three axes. You do not get all three at the top.
| Setting | Freshness | Query flexibility | Cost per query | When |
|---|---|---|---|---|
| Raw stream ingest, no rollup | 5–30 s | Full | Highest | Ops/debug dashboards with tens of QPS |
| Stream ingest + minute rollup | 5–30 s | Fixed dimensions | Low | Ad pacing, merchant "today" tiles |
| Stream head + hourly batch tail (hybrid) | 5–30 s head, hours tail | Fixed in tail | Lowest at scale | User-facing products over months of history |
| Batch ingest only, every 15–60 min | 15–60 min | Fixed | Lowest; simplest ops | Reports that nobody reads within the hour |
| Warehouse, no OLAP store | 5 min to hours | Full SQL, joins | Seconds and $/TB scanned | Internal analysts |
The cost formula to say out loud:
cpu_seconds_per_sec ≈ QPS × bytes_scanned_per_query / scan_rate_per_core
Merchant dashboard, no rollup:
5,000 QPS × 200 MB (7 days of one large tenant, 4 columns, compressed) / 1 GB/s/core
≈ 1,000 cores busy
Same dashboard on a 1-minute rollup:
5,000 QPS × 4 MB / 1 GB/s/core ≈ 20 cores
A 50× difference in cores is a 50× difference in the bill, and it came entirely from a modeling decision. The worst common query, the large tenant's 7-day view, sets the cluster size, not the median one.
🎯 Staff Insight: "I size OLAP by the 99th-percentile tenant's query times peak QPS. The median tenant reads a few kilobytes and is free. The top 1% of tenants by data volume are the cluster."
Anti-Patterns — What Kills Real-Time OLAP Deployments#
1. Using It as the System of Record#
The OLAP table becomes the only place an order's final status lives, because it was convenient. Then a schema change requires re-ingestion and there is no source to re-ingest from. Fix: Kafka with long enough retention plus raw events in object storage are the record; the store must be rebuildable within a day.
2. Row-at-a-Time Inserts into ClickHouse#
A service writes one INSERT per event at 20K events/s. Parts pile up faster than merges can combine them, queries slow as they open thousands of parts, and inserts start failing with "too many parts". Fix: batch at the client (10K–100K rows or 1 s), use the Kafka engine, or async inserts with wait_for_async_insert=1.
3. The Sort Key Picked From the Event Schema#
ORDER BY (event_id) or ORDER BY (event_time) on a multi-tenant table. Every tenant query scans the whole time range for all tenants. Fix: dominant filter first, usually tenant_id, then a coarse time bucket. Changing it later means rewriting the table.
4. Joins on the Hot Path#
User-facing queries join a 10B-row fact table with a 50M-row dimension table at request time. Distributed joins shuffle data between servers and blow up memory. Fix: denormalize at ingest (Flink enrichment or lookup tables), keep small dimension tables replicated to every node, and leave large joins to the warehouse.
5. Exact Distinct Counts Everywhere#
COUNT(DISTINCT user_id) over 30 days on every dashboard load. Each server ships its full set of IDs to the broker. Fix: sketches by default; exact counts only for the one billing or compliance report, computed in batch.
6. Unbounded Multi-Tenancy Without Quotas#
One enterprise tenant with 40% of the data shares servers with 2M small tenants. Their 90-day export query saturates every server in the fan-out, and p99 for everyone else goes from 80 ms to 4 s. Fix: per-tenant query quotas and timeouts, separate tenant tiers or tables for the largest tenants, and result caching for the dashboard's fixed queries.
7. Too Many Tiny Segments or Partitions#
Partitioning by day × tenant, or sealing real-time segments every minute. Queries fan out to tens of thousands of segments and spend their time in scheduling, not scanning. Fix: monthly partitions in ClickHouse; 300–700 MB segments in Druid; compaction jobs owned and alerted on.
8. "Real Time" Without a Lag SLO#
Nobody defines freshness, so nobody notices when ingest lag drifts from 10 s to 40 minutes during a traffic spike. The dashboard looks calm while the incident burns. Fix: publish max(event_time) per table as a metric and show it on the dashboard; alert on Kafka consumer lag in seconds, not messages.
The Technology Landscape — Head-to-Head Comparison#
| Dimension | ClickHouse | Apache Druid | Apache Pinot | Cloud Warehouse (BigQuery, Snowflake, Redshift) |
|---|---|---|---|---|
| Designed for | Fast SQL analytics on one table at a time | Interactive slice-and-dice over event streams | User-facing analytics at very high QPS | Ad-hoc analysis, ELT, joins across the business |
| Typical p99 | 10 ms – seconds, depends on modeling | 50 ms – 1 s | 10–100 ms for tuned queries | Seconds to minutes |
| Concurrency comfort | Hundreds of QPS per cluster before careful tuning | Thousands | Thousands to tens of thousands | Tens, with queueing and per-query billing |
| Stream freshness | Seconds (Kafka engine or batched inserts) | Seconds (Kafka indexing service) | Seconds (direct partition consumption) | Minutes typically; streaming inserts exist at extra cost |
| Joins | Good: hash joins, dictionaries; large joins need care | Limited; lookups and broadcast joins | Limited; lookup joins, multi-stage engine | Full |
| Mutability | Eventual via merges; mutations are heavy | Interval re-ingest | Upsert and dedup on realtime tables | Full DML |
| Ops footprint | Smallest: one binary + Keeper | Largest: 5–6 process types + ZooKeeper + metadata DB + deep storage | Medium: controller, broker, server, minion + ZooKeeper + deep store | None (managed) |
| Cost shape | Hardware you run | Hardware you run | Hardware you run | $ per TB scanned or per compute-hour |
Staff reading of the table: choose ClickHouse when one team wants SQL flexibility and the QPS is in the hundreds. Choose Pinot when the product is "every user sees their own analytics" with thousands of QPS, upserts or tenant-scoped indexes. Choose Druid when time-series slicing with ingest-time rollup is the core and the organisation can run its process zoo. Choose the warehouse you already have when the users are analysts.
Which Store? A Decision Tree#
Patterns#
Pattern 1: Kafka → OLAP Direct (Kappa Serving)#
Producers write to Kafka keyed by tenant; the OLAP store consumes directly; dashboards query the store. Use when events are already clean and denormalized, and freshness under 30 s matters. Watch: no place to fix bad data except re-ingest; schema changes ripple straight into the table.
Pattern 2: Stream Processor in Front (Enrich, Dedup, Re-key)#
Flink sits between raw topics and the OLAP topic: joins dimension data, drops duplicates by event ID within a window, re-partitions by the store's primary key. Use when upserts need key partitioning, events need enrichment to avoid query-time joins, or exactly-once counting matters. This is the shape in Ad Click Aggregation and Stream Processing; see Flink for the checkpointing side.
Pattern 3: Hybrid Tables — Stream Head, Batch Tail#
Pinot's hybrid table (realtime + offline with a time boundary) and Druid's stream-plus-batch re-ingestion share one idea: the last 1–3 days come from Kafka with best-effort rollup; older data is rebuilt nightly from object storage with perfect rollup, corrections applied and larger segments. Use when history is long and correctness of old data matters. Watch: the boundary. A query spanning it reads two sources; the batch job that fails silently leaves a hole that only the realtime side was hiding.
Pattern 4: Layered Rollups (Raw → Minute → Hour → Day)#
Multiple tables at different grains with different TTLs: raw for 7–14 days, minute for 90 days, hour for 13 months, day forever. The API routes each query to the coarsest table that can answer it. Use when queries span very different time ranges. Watch: route logic belongs in one query service, not every client, or a new dashboard will scan raw data for a 1-year chart.
Pattern 5: Query Service With a Result Cache#
User-facing dashboards issue the same 5 queries per tenant. A thin API layer owns the query templates, adds tenant filters (never trust the client), enforces timeouts and per-tenant quotas, and caches results for 10–60 s keyed by (tenant, query, time bucket). Use when always, for user-facing analytics. The cache absorbs a page refresh storm; the template list doubles as the index design document.
Pattern 6: Tenant Tiering#
The largest 0.1% of tenants get their own tables or server tenants (Pinot supports tagging servers for specific tables; ClickHouse can route big tenants to a separate cluster). Use when tenant size is power-law distributed, which in B2B analytics it nearly always is. See Hot Keys for the general version of this problem.
Scaling#
The Numbers to Size With#
| Quantity | Planning Number (state assumptions) |
|---|---|
| Compressed bytes per event (10–20 columns, sorted well) | 15–60 bytes |
| Scan rate per core, simple aggregations on compressed columns | ~0.5–2 GB/s of uncompressed data |
| Kafka partition → consumer | 1:1; ~10–50K events/s per consumer is a safe planning range |
| Druid segment target | 300–700 MB, ~5M rows |
| ClickHouse insert batch | ≥ 1,000 rows, ideally 10K–100K |
| ClickHouse granule | 8,192 rows default |
| Star-tree leaf threshold | 10,000 records default |
| Pinot upsert key map | ~(key bytes + 24) bytes per live key |
| Freshness, streaming ingest | 5–30 s typical |
Sizing Walkthrough#
Ad events: 2M events/s peak, 1M average → ~86B events/day.
Raw storage: 86B × 40 B ≈ 3.4 TB/day compressed → 14 days ≈ 48 TB × RF 2 ≈ 96 TB
Minute rollup: ratio 30× ≈ 115 GB/day → 13 months ≈ 45 TB × RF 2 ≈ 90 TB
Ingest: 2M/s ÷ 25K per consumer ≈ 80 Kafka partitions minimum → choose 128
Query: 3,000 QPS, worst common query scans 50 MB of rollup
3,000 × 50 MB / 1 GB/s/core ≈ 150 cores busy → ×2 headroom ≈ 300 cores
Servers: ~16 servers × 32 cores, NVMe for the hot 30 days, object storage for the rest
The two numbers that set the shape: Kafka partition count (ingest parallelism, fixed early) and worst common query × QPS (CPU). Storage is the cheapest line.
What Breaks First as You Grow#
| Growth Axis | First Thing to Break | Move |
|---|---|---|
| Event rate | Consumers fall behind; part or segment count explodes | More partitions (planned up front), bigger batches, compaction capacity |
| Query QPS | Broker merge CPU and fan-out to many servers | Partition-aware routing so a tenant hits 1–2 servers; replica groups; result cache |
| Data history | Segment count and metadata; cold queries hit object storage | Tiered storage; coarser rollups for old data; fewer, bigger segments |
| Tenants | One giant tenant sets everyone's p99 | Tenant tiering; quotas; dedicated tables |
| Dimensions | Rollup ratio collapses; dictionaries bloat | Dimension admission; separate tables for high-cardinality questions |
Multi-Region#
Real-time OLAP is almost always regional and rebuildable: each region consumes its own (or a mirrored) Kafka stream into its own cluster. Do not stretch one cluster across regions: replication and fan-out latency between regions destroys the p99 that justified the store. For global dashboards, either mirror topics to every region and accept seconds of skew, or route the tenant to its home region. The DR plan is re-ingest from Kafka plus object storage, which means Kafka retention must cover the time to rebuild.
Failure Modes & Recovery#
1. Ingest Lag Spiral During a Traffic Spike#
t=0 Promo launches; event rate 3x normal
t=+2min Consumers at 100% CPU building in-memory segments; lag 90 s and climbing
t=+6min Segment seal + upload slows under disk contention; lag 8 min
t=+10min Dashboards show calm numbers from 8 minutes ago; ops misses the incident
- Root cause: ingest capacity sized for average, not peak; freshness not visible to users.
- Detection:
kafka_consumer_lag_secondsper partition;ingest.max_event_time_ageper table; segment seal duration. - Mitigation: add consumers only if partitions allow; shed optional columns or enrichments; show a "data delayed" banner driven by event-time age.
- Prevention: size consumers for 2–3× peak; partition count chosen for peak; lag SLO with a page at 2 minutes.
- Owner: data platform on-call; product owns the banner behaviour.
2. Too Many Parts / Segment Explosion#
- Symptom: ClickHouse inserts fail with "too many parts"; or Druid/Pinot queries slow as segment count passes tens of thousands.
- Root cause: small inserts, overly fine partitioning, or a stalled compaction job.
- Detection:
active_parts_per_partition,merges_in_progress,segment_countper table, compaction task failures. - Fix: batch upstream; pause the offending writer; run catch-up compaction on dedicated capacity.
- Prevention: an insert-size guard in the ingest client; partition key review in schema change process; compaction owned and alerted.
3. The Noisy Tenant Query#
- Symptom: p99 for all tenants jumps from 80 ms to 3 s; CPU pinned on every server.
- Root cause: one large tenant (or an internal export) runs an unbounded 90-day raw scan with a high-cardinality group-by.
- Detection: per-tenant
query_cpu_secondsandbytes_scanned; slow-query log by tenant. - Fix: kill the query; apply a per-tenant quota; route that tenant to its tier.
- Prevention: query templates in a service layer; per-query timeout (e.g. 5 s for user-facing);
max_bytes_to_read-style guards; tenant tiering. - Owner: the query service team owns limits; account management owns the conversation with the tenant.
4. Silent Double Counting After a Replay#
- Symptom: revenue tile 1.8× higher for one hour after a consumer redeploy; finance asks why it disagrees with the warehouse.
- Root cause: at-least-once ingest replayed a batch; rollup MV counted it twice;
ReplacingMergeTreeon raw rows does not fix an aggregate already written. - Detection: a reconciliation job comparing hourly totals with the warehouse or the stream processor's counts; alert at > 0.5% drift.
- Fix: re-aggregate the affected hours from deduplicated raw data and replace them.
- Prevention: dedup by event ID upstream (stream processor) or exact upsert/dedup tables; idempotent, versioned batch replacement.
- Owner: data platform; finance signs off on the tolerance.
5. Upsert Memory Exhaustion#
- Symptom: Pinot servers hit GC storms or OOM; consuming segments stop; lag climbs on every upsert table on those servers.
- Root cause: primary-key map grows without bound because keys never expire; or a producer bug floods new keys.
- Detection: heap usage;
upsert_primary_keys_countper partition; GC pause time. - Fix: add servers for the replica group; enable metadata TTL; quarantine the bad producer.
- Prevention: key-count budget per table at design review; TTL on upsert metadata; offline compaction of cold data.
Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Ingest lag spiral | consumer lag seconds; event-time age | All dashboards on the table | Scale consumers, shed enrichments, banner | Data platform |
| Part/segment explosion | part and segment counts | Whole table, then cluster | Batch, compact, pause writer | Data platform |
| Noisy tenant | per-tenant CPU and bytes scanned | Every tenant on shared servers | Kill, quota, tier | Query service team |
| Double counting | reconciliation drift | Business decisions, invoices if misused | Re-aggregate the interval | Data platform + finance |
| Upsert OOM | heap, key count | Every upsert table on the replica group | TTL, scale, quarantine | Data platform |
| Deep storage outage | upload failures | New segments stay on servers; disks fill | Extend local retention; pause compaction | Data platform + storage |
When to Use vs. Alternatives#
| Situation | Pick | Why |
|---|---|---|
| End-user analytics, thousands of QPS, fixed queries, seconds freshness | Pinot or Druid (or ClickHouse with MVs + cache) | Built for concurrency and predictable latency |
| Internal product or log analytics, ad-hoc SQL, dozen to hundreds of QPS | ClickHouse | SQL flexibility, joins, smallest ops footprint |
| Analysts, cross-domain joins, minutes freshness fine | The warehouse you already pay for | No new system to run |
| Infra metrics with labels and alerting | TSDB (Time-Series Databases) | PromQL, alerting, rate/histogram semantics |
| Full-text search over logs | Elasticsearch | Inverted index on text, relevance |
| Under ~100M–1B rows, low QPS, one team | PostgreSQL with rollup tables | Free in operational terms; you already run it |
| Money, invoices, payouts | Warehouse of record or OLTP with reconciliation | Exactness and audit, not milliseconds |
When NOT to Use a Real-Time OLAP Store#
- Instead of the warehouse. If the users are 30 analysts and freshness is "this morning", a real-time OLAP cluster is an extra on-call rotation that buys nothing. The warehouse charges per query; the OLAP store charges per hour whether anyone queries it or not. Add the OLAP store when concurrency or latency, not convenience, demands it.
- Instead of a TSDB. For infrastructure metrics with alerting, a TSDB gives you rate(), histograms, recording rules and an alerting pipeline that fails independently. An OLAP store is the right home for the high-cardinality questions a TSDB refuses (per-customer, per-URL), not for the pager. See Metrics & Alerting Platform.
- Instead of Elasticsearch. If the core query is "find log lines containing this phrase", you need an inverted index on text. ClickHouse and Pinot have text indexes, but relevance, analyzers and highlighting are Elasticsearch's job. Conversely, Elasticsearch is an expensive aggregation engine; heavy dashboards over logs often belong in ClickHouse.
- Instead of Postgres. Under ~1B rows and a few dozen QPS, Postgres with rollup tables refreshed every minute, or a columnar extension, is cheaper than any new cluster. Moving off it is a two-way door; adopting a distributed OLAP engine is closer to one-way.
- As the system of record or for transactions. No multi-row transactions, eventual dedup, weak updates. If correctness matters, the record lives elsewhere and the OLAP store is a view of it.
- For point lookups by primary key at high QPS. A key-value store answers
GET order 123in 1 ms at a fraction of the cost.
🎯 Staff Insight: "The question that kills most real-time OLAP proposals is: who looks at this number within a minute of the event, and what do they do differently? If the honest answer is nobody, a 15-minute batch into the warehouse is the design."
Operational Concerns#
What the On-Call Actually Does#
- Watches freshness first. Event-time age per table is the user-visible SLI; consumer lag is the cause.
- Runs and repairs compaction. Failed compaction is invisible for days, then shows up as slow queries and full disks.
- Manages schema changes. Adding a column is cheap; changing a sort key, partition key or rollup dimension is a re-ingest. Those go through review with a backfill plan.
- Handles backfills. Corrections arrive as "re-ingest last Tuesday": a batch job that writes new segments or partitions and atomically swaps them, with a reconciliation check before and after.
- Polices tenants. Top-N tenants by CPU and bytes scanned every week; quota changes are a product conversation, not an infra secret.
Key Metrics & Alerts#
| Metric | Alert Threshold (example) | Why |
|---|---|---|
ingest.event_time_age_seconds | > 120 s for 5 min | The freshness promise |
kafka_consumer_lag_seconds | > 60 s and rising | Leading indicator of the above |
query.latency_p99_ms per query template | > 2× baseline for 10 min | User-facing SLO |
query.bytes_scanned per tenant | top tenant > 30% of cluster | Noisy-neighbour early warning |
parts_per_partition / segment_count | ClickHouse > 300 active parts in a partition | Precursor to insert failures |
rollup_ratio | drops > 30% week over week | A new dimension or a cardinality bug |
reconciliation_drift_pct | > 0.5% hourly | Duplicates or loss |
upsert_primary_keys per server | > 80% of budget | OOM ahead |
Schema Evolution and Backfills#
Every engine handles adding a nullable column online. Everything else is a rebuild: new sort key, new partitioning, new rollup grain, new star-tree split order. The rebuild pattern is the same in all three: build the new table from raw events in object storage, dual-read and compare for a week, switch the query service, retire the old table. Budget it in days of compute and an engineer-week, and keep raw events long enough that it is possible at all.
Interview Application — Staff-Level Plays#
Which Case Studies Use Real-Time OLAP#
| Case Study | How Real-Time OLAP Is Used | Key Pattern |
|---|---|---|
| Ad Click Aggregation | Druid or Pinot serving per-campaign counts to advertiser dashboards and pacing | Dedup upstream, rollup at ingest, warehouse reconciles billing |
| Stream Processing | Serving store downstream of Flink windows, upserted by (campaign, window_start) | Upsert key = window key; replay overwrites, never adds |
| Metrics & Alerting Platform | High-cardinality questions (per customer, per endpoint) that a TSDB refuses | TSDB for alerting; OLAP for exploration |
| Leaderboard | "Top products this quarter" style analytics boards | Redis for live ranks; OLAP for ad-hoc slicing |
| Database Selection | The derived analytical store fed by CDC | Never run scans on the OLTP primary |
| URL Shortener | Per-link click analytics behind the redirect path | Async event pipeline; analytics never on the hot path |
| Search Indexing | Where aggregation-heavy queries go instead of the search cluster | Search clusters are expensive analytics engines |
Related reading: Kafka for partitioning and retention, Batch & Stream Pipelines for the batch and stream layers that feed the store.
Every System Design Question Has a Real-Time OLAP Moment#
- URL shortener: "Click analytics go redirect → Kafka → Pinot, sorted by
link_id. The redirect never waits on analytics, and the creator's chart is 30 seconds behind, which I'll show on the page." - Video streaming: "Player heartbeats for QoE go to a real-time OLAP store with minute rollups by device, CDN and region, so a canary release can be compared against control within minutes. See Video Streaming."
- Payments: "Merchant sales dashboards read from OLAP; payouts never do. The payout ledger is the record, and the dashboard says 'estimated' until settlement."
- Feature flags / experiments: "Exposure and conversion events land in OLAP for live guardrail metrics; the final experiment readout runs in the warehouse with exact joins."
What Interviewers Probe#
| After You Say... | They Will Ask... | What They're Evaluating |
|---|---|---|
| "Store events in ClickHouse" | "What's the sort key? How many inserts per second?" | Layout from queries; small-insert failure mode |
| "Real-time dashboards" | "How fresh, and how does the user know when it isn't?" | Lag SLO and visibility |
| "Count unique users" | "Exactly? Over 90 days? Per tenant?" | Sketches and their error |
| "Kafka is at-least-once, so..." | "What happens to your rollup on a replay?" | Dedup placement before aggregation |
| "Pinot with upserts" | "How much memory, and what's the partitioning requirement?" | Knows the price of exactness |
| "Multi-tenant" | "Your biggest tenant has 40% of the data. Now what?" | Tenant tiering and quotas |
Common Interview Mistakes#
| What Candidates Say | What Interviewers Hear | What Staff Engineers Say |
|---|---|---|
| "ClickHouse is fast, it'll handle it" | No model of cost per query | "Bytes scanned per query times QPS: 3,000 × 50 MB is ~150 cores." |
| "Store every field, decide later" | Rollup ratio of 1; scans grow forever | "Raw for 14 days in a narrow table, minute rollup for 13 months." |
| "Exactly-once from Kafka to the store" | Hasn't looked at merge-time dedup | "At-least-once into the store; dedup by event ID before the rollup." |
| "Join with the users table at query time" | Distributed join on the hot path | "Enrich in Flink at ingest; the store sees one wide row." |
| "Use OLAP for the invoices too" | Money on an eventually-deduped store | "Invoices come from the warehouse of record; OLAP is labelled estimated." |
| "We'll use Druid" for 20 analysts | Five process types for a warehouse job | "That's the warehouse. Real-time OLAP when concurrency demands it." |
L5 vs L6 vs L7 Responses#
| Scenario | L5 Answer | L6 / Staff Answer | L7 / Principal Answer |
|---|---|---|---|
| "Build merchant analytics for 2M sellers" | Events into ClickHouse, dashboard queries it | Query shapes first; tenant-first sort key; minute rollup; HLL for uniques; query service with templates, quotas, cache; lag banner | One analytics platform with a dimension-admission process and per-product chargeback; decides which numbers are certified and where they come from |
| "Dashboards slowed from 80 ms to 3 s" | Add servers | Find the tenant and the query by bytes scanned; quota it; tier the top tenants; fix the template | Introduces per-tenant cost budgets and makes the largest customers' tier a priced product |
| "Numbers disagree with finance" | Re-run the query | Reconciliation job, find the replay, move dedup before the rollup | Publishes which tables are certified, the tolerance finance accepted, and who signs off on corrections |
| "Should we adopt Pinot?" | Yes, it's what LinkedIn uses | Only if QPS > ~1K or upserts needed; otherwise ClickHouse or the warehouse | Models 3-year cost including the platform team, names the exit path, and sets the trigger to revisit |
The Staff Real-Time OLAP Checklist#
- Query population: "Who queries, how many at once, at what p99? Analysts at 20 QPS or users at 5,000?"
- Query shapes: "Five templates: filters, group-bys, time ranges. The sort key is
(tenant_id, time)." - Grain and retention: "Raw 14 days, minute rollup 13 months; product signed off on no single-order drill-down after two weeks."
- Correctness: "At-least-once ingest, dedup by event ID upstream of the rollup, sketches for distinct counts, warehouse for anything billed."
- Freshness: "Kafka-direct ingest, 30-second target, event-time age on the dashboard, page at 2 minutes."
- Cost and isolation: "Bytes scanned × QPS sizes the cluster; per-tenant quotas and a tier for the top 0.1% of tenants."
🎯 Staff Insight: Don't use a real-time OLAP store as the system of record, for invoices, for high-QPS point lookups, as a TSDB for alerting, for full-text search, or when nobody acts on the number within the hour. The strongest signal is saying what the dashboard shows when ingest is 10 minutes behind.
Evaluation Rubric#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Modeling | Wide table of events | Sort key, grain, rollup and indexes derived from named query shapes | Dimension admission policy; schema changes reviewed with a backfill budget |
| Correctness | Assumes counts are right | Places dedup before aggregation; knows merge-time dedup is eventual; uses sketches knowingly | Certified vs directional data, reconciliation tolerance signed off by finance |
| Performance | "It's columnar, it's fast" | Bytes scanned × QPS; pruning; tenant tiering | $ per 1,000 queries per product; chargeback |
| Operations | "Monitor the cluster" | Lag SLO, compaction, part counts, noisy tenants | Rebuild-from-source drill; regional DR via re-ingest; platform staffing |
| Choice | The engine they know | Picks by concurrency, upserts, joins, ops footprint | Build vs buy, exit path, when to keep the warehouse |
Strong hire signals
| Signal | What It Sounds Like |
|---|---|
| Query-first modeling | "Give me the five dashboard queries and I'll give you the sort key." |
| Prices queries | "The top tenant's 7-day view is 50 MB; at 3,000 QPS that's the cluster." |
| Correctness placement | "Replays hit the rollup before any merge dedups them, so dedup goes upstream." |
| Freshness as a product feature | "The page shows data-as-of; lag over two minutes pages us." |
| Knows when not to | "Twenty analysts and hourly freshness is the warehouse." |
Lean no-hire signals
| Signal | Why It Misses the Bar |
|---|---|
| No sort key or grain discussed | Will scan everything and buy hardware to compensate |
| Exact distinct counts everywhere | Unbounded memory and fan-out cost |
| OLAP as the source of truth | No rebuild path; correctness unowned |
| No tenant isolation in a multi-tenant product | One customer sets everyone's p99 |
Common false positives
- Knowing star-tree internals ≠ modeling judgment. Ask which queries the tree serves and what it costs to change the split order.
- "We ingested 1M events/s" ≠ serving. Ask about p99 at peak QPS and the noisy-tenant story.
- Engine benchmark fluency ≠ design. Vendor benchmarks measure scans on one table; production pain is concurrency, freshness and correctness.
Beyond Staff: The Principal View#
Why L7 Sees This Problem Differently#
At Staff level a real-time OLAP store is a serving system you model correctly. At Principal level it is a shared bill with an unbounded number of authors. Every product team wants its own dashboard, every dashboard wants one more dimension, and every dimension multiplies rows, memory and query cost for everyone on the cluster. The L7 question is not "Pinot or ClickHouse?" but "How many analytics engines does this company run, who may add a dimension or a tenant tier, which numbers are certified, and what does a query cost the team that issues it?"
🧭 Principal Move: "Before choosing an engine, I want two lists signed off: the certified metrics, which come from the warehouse of record and appear on invoices and board decks, and the operational metrics, which come from the real-time store and are labelled with their freshness and error. Most of the cost and most of the incidents in real-time analytics come from blurring those two lists."
The Org-Level Fault Line#
One shared real-time analytics platform vs a store per product team.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| One shared platform (one engine, multi-tenant clusters) | One expert team; shared ingest and tooling; consistent metric definitions | One team's dimension or query hurts everyone; platform becomes a bottleneck for schema changes | Platform team and every product during noisy-neighbour incidents |
| Shared platform with tiers and chargeback | Isolation for heavy users; cost visible to the teams creating it | Needs metering, quotas and an admission process | Platform headcount for the control plane |
| Store per product team | Autonomy; fast iteration | 3 engines, 6 clusters, inconsistent definitions of "active user", nobody to page at 3am | Finance through duplication; data consumers through conflicting numbers |
| Warehouse only, no real-time store | Lowest ops; single source | Cannot serve user-facing QPS or seconds-fresh dashboards | Product: features that need freshness don't ship |
The Principal default: one engine family on a platform-owned, tiered deployment, a metrics layer that defines each certified metric once, per-team chargeback on CPU and storage, and a written exception path for teams whose needs genuinely differ (for example, one ClickHouse cluster for log analytics alongside Pinot for user-facing products). Two engines is a deliberate choice; four is drift.
🧭 Principal Insight: "The expensive thing isn't the cluster, it's an unpriced dimension. If adding
product_idto the shared rollup is free to the team that asks, it will be added, and everyone's bill goes up 5×."
Cost Model#
Assumptions: storage-optimised server 32 vCPU / 256 GB / NVMe at ~$2K/month; object storage ~$23/TB-month; Kafka share priced separately; loaded engineer $250K/year ($21K/month). Directional only.
| Scale | Workload | Servers + storage/month | Kafka + pipelines/month | Platform headcount | Total/month |
|---|---|---|---|---|---|
| Startup | 50M events/day, 200 QPS, one product | 3 nodes ~$6K | ~$2K | 0.5 FTE (~$10K) | ~$18K, versus ~$3–8K for warehouse plus cached queries |
| Growth | 5B events/day, 3K QPS, 5 products | 16 nodes + 100 TB object ~$35K | ~$10K | 2 FTE (~$42K) | ~$87K |
| Enterprise | 100B events/day, 20K QPS, 30 products | 120 nodes + 2 PB object ~$290K | ~$70K | 7 FTE (~$150K) | ~$510K |
At startup scale the engineer is the dominant cost, which is the quantitative case for staying on the warehouse plus a cache until concurrency forces the move. At enterprise scale, rollup ratio and the top tenants' queries are the levers: raising the average rollup ratio from 5× to 15× roughly halves the server line, more than any hardware negotiation.
The 3-Year Evolution Path#
One-Way Doors vs Two-Way Doors#
| Decision | Reversibility | Cost to Reverse |
|---|---|---|
| Kafka partition count for upsert tables | One-way for that table | New topic, new table, full re-ingest |
| Rollup grain with no raw events kept | One-way | Lost detail is gone; only future data can be finer |
| Sort key / partition key / star-tree split order | One-way-ish | Rebuild from raw events: days of compute, an engineer-week |
| Exposing OLAP numbers in a customer-facing API contract | One-way-ish | Customers build on its semantics and freshness |
| Engine choice (ClickHouse vs Pinot vs Druid) | One-way-ish | Quarters: ingestion specs, query templates, operations runbooks |
| Retention per tier | Two-way going shorter, one-way for data already deleted | Storage cost only, until something is deleted |
| Adding a nullable column, a new MV, a new index | Two-way | Backfill time |
| Result cache TTL, query quotas | Two-way | Config change |
The Standard I'd Write#
RFC-DATA-011: Real-Time Analytics Tables
Scope: Every table in the shared real-time analytics platform serving
dashboards, product features or external APIs.
MUST
1. Declare the query templates the table serves, its sort/partition key,
rollup grain, retention per tier and a freshness SLO.
2. Be rebuildable from Kafka and object storage; raw events retained
>= 14 days or the rebuild window, whichever is longer.
3. Deduplicate by event ID before any aggregation; state the delivery
guarantee and expected error in the table description.
4. Expose data-as-of (max event time) to every consumer.
5. Pass a cardinality review for each new dimension: observed distinct
values, effect on rollup ratio, cost estimate, owning team.
6. Never serve as the source for invoices, payouts or regulatory reports.
SHOULD
7. Use sketches for distinct counts; exact only by exception.
8. Route through the query service (templates, quotas, timeouts, cache).
9. Place tenants above 5% of table volume in a dedicated tier.
Exceptions: Filed with the data platform team; approved by the platform lead
and the consuming product's director; reviewed quarterly.
Rollout: shadow (report violations) -> warn -> enforce on new tables ->
enforce on existing tables within two quarters.
Success metrics: freshness SLO met >= 99.5% of minutes; zero certified
metrics served from OLAP; reconciliation drift < 0.5%; cost per 1,000
queries flat or falling while QPS grows.
What I'd Tell the VP#
"Our customers now expect to see their own numbers within a minute, and the warehouse can't serve that many people that fast, so we need a real-time analytics system. It will cost roughly $90K a month at our current size, including two engineers, and it grows with how many questions we let teams ask, not with how much data we keep. I'm proposing one shared platform rather than one per team, with each team paying for what its dashboards use, so the bill stays visible. The numbers it shows will be labelled as live estimates; anything we invoice or report to the board still comes from the warehouse, so a bug here can't turn into a billing error."
Principal Interview Signals#
| Signal | What It Sounds Like |
|---|---|
| Separates certified from operational data | "Invoices from the warehouse, dashboards from OLAP, and the page says which." |
| Prices dimensions | "Each new dimension gets a cardinality review and a line on the requesting team's bill." |
| Limits engine sprawl | "Two engines on purpose, with a written reason. A fourth is drift." |
| Plans rebuildability | "If we can't rebuild a table from raw in a day, it's a database we didn't design." |
| Sets adoption triggers | "Warehouse until user-facing QPS passes about a thousand or freshness drops under a minute." |
Staff answers that L7 interviewers find insufficient:
- "We'll add quotas per tenant" without saying who sets them, how customers are told, or how a top tenant buys more.
- "Reconcile against the warehouse" without the tolerance, the owner of the drift alert, or what happens to the dashboard when it fires.
- "Pick Pinot for scale" without the platform headcount, the exit path, or the trigger that justified leaving the warehouse.
How Real Companies Built It#
LinkedIn — Pinot#
LinkedIn built Pinot to serve analytics directly to members and customers, the "Who Viewed My Profile" kind of feature, rather than only to internal analysts. Its SIGMOD 2018 paper, Pinot: Realtime OLAP for 530 Million Users, describes a single production system that serves tens of thousands of analytical queries per second, ingests in near real time from streaming sources, and argues that relational databases and key-value stores fall apart at that combination of high ingest rate and high analytical query rate at low latency (Pinot paper, SIGMOD 2018).
Staff insight: Pinot's design follows from the query population: millions of users each asking a small, tenant-scoped question. That is why it invests in pruning, partition-aware routing and the star-tree. In an interview, name the population before the engine.
Cloudflare — ClickHouse for HTTP Analytics#
Cloudflare rebuilt its HTTP analytics pipeline on ClickHouse, replacing a pipeline built around a single PostgreSQL rollup database and a Citus cluster. It described processing about 6M HTTP requests per second on average, with peaks of 8M, on a 36-node ClickHouse cluster with 3× replication, using SummingMergeTree and AggregatingMergeTree tables fed by materialized views, and reported that lowering index granularity from 8,192 to 32 on the aggregated tables cut query latency by about half and raised throughput about 3× (Cloudflare blog).
Staff insight: The win came from modeling, not hardware: pre-aggregated tables via materialized views, and an index granularity tuned for small, tenant-scoped reads instead of the scan-oriented default. The default settings serve the analyst; serving a dashboard means re-tuning them.
Netflix — Druid for Playback Quality#
Netflix uses Apache Druid to watch the quality of playback across devices in real time, using logs from playback devices to derive quality measures tagged by device characteristics. It described data arriving at over 2 million events per second and more than 115 billion rows per day, and uses the metrics to compare a new software version enabled for a subset of users against the existing one, so that a regression can abort the rollout (Netflix Technology Blog).
Staff insight: This is the operational-monitoring intent: high-cardinality slicing (device type, version) with seconds of freshness, where approximate is fine and stale is not. The OLAP store sits inside a release-safety loop, so its lag SLO is part of the deployment system's correctness.
Druid's Origin and Uber's Ad Pipeline#
Druid itself came out of ad-tech analytics at Metamarkets, and Uber's exactly-once ad pipeline uses Pinot upserts as its serving layer. Both are covered, with sources, in Ad Click Aggregation; the Druid design is described in its SIGMOD 2014 paper (Druid paper, SIGMOD 2014).
Practice Drill#
Prompt: "You run a marketplace with 3M sellers. Seller analytics (sales today, top products, traffic sources) are served from Postgres materialized views refreshed every 30 minutes. Sellers complain the numbers are stale, the refresh now takes 25 minutes, and product wants 'live' analytics plus a public API. Design the replacement."
Staff Answer
Start with the query population: 3M sellers, maybe 150K active in the busiest hour, five dashboard tiles each, so peak ~3–5K QPS of tenant-scoped, time-bounded queries with a 300 ms page budget. That rules out the warehouse for serving and points at Pinot (or ClickHouse with materialized views and a cache if the team prefers SQL and QPS stays under ~1K). Pipeline: order and page-view events to Kafka keyed by seller_id, 64 partitions sized for 3× peak; a Flink job dedups by event ID, enriches with product category, and handles order status changes. Tables: an upsert table for orders keyed by order_id (refunds show within a minute; key window 30 days, ~(8+24) bytes per key budgeted), and an append table for traffic with a minute-grain star-tree over (source, device) and SUM metrics. Sort/partition by seller_id, so a seller's query touches one partition. Unique visitors via HLL. A query service owns the five templates, injects seller_id, enforces a 2 s timeout and per-seller quotas, and caches for 15 s. The largest 0.1% of sellers go to a dedicated tier. Freshness target 30 s, data_as_of on the page, page the on-call at 2 minutes of lag. The public API exposes the same templates with documented freshness and "estimated" semantics; payouts stay on the ledger. Migration: dual-write for two weeks, compare OLAP totals with Postgres hourly, cut over per tile, then retire the views.
Why this is L6:
- Derives the engine from QPS and tenant-scoped query shapes, not preference.
- Places dedup before aggregation and prices the upsert memory.
- Isolates large sellers and protects p99 with templates, quotas and a cache.
- Makes freshness visible and keeps money on the ledger.
What L7 adds:
- Treats the public API as a one-way door: freezes metric definitions and freshness semantics in a contract before launch.
- Prices it: ~$40–60K/month plus 2 engineers, against seller churn attributed to stale analytics, and sets the review trigger for adding a second product on the same platform.
- Writes the dimension-admission and certified-vs-operational policy so the next five teams inherit it instead of inventing their own stores.
Quick Reference Card#
Model: pre-paid answers to known questions; sort key = index, rollup = cost
Pick: analysts -> warehouse; ad-hoc SQL under ~1K QPS -> ClickHouse;
users at 1K-50K QPS, upserts -> Pinot; time-series rollup -> Druid/Pinot
Capacity unit: bytes scanned per query x QPS (storage is the cheap line)
Sort key: (tenant_id, coarse time, next filter); partition by month, never by tenant
Freshness: Kafka-direct ingest 5-30 s; publish data-as-of; page at ~2 min lag
Ingest: 1 Kafka partition = 1 consumer; partition count chosen for 3x peak
ClickHouse: batches >= 1,000 rows (10K-100K ideal); granule 8,192 rows;
ReplacingMergeTree dedups only on merge -> FINAL or argMax
Druid: segments 300-700 MB (~5M rows); streaming rollup is best-effort
Pinot: star-tree maxLeafRecords 10,000 default; all aggs in functionColumnPairs;
upsert needs stream partitioned by PK, ~(key+24) B heap per key,
strictReplicaGroup routing, no star-tree
Rollup ratio: events / rows; below ~5x the dimension set is too wide
Distinct counts: HLL / Theta sketches; exact only in batch
Correctness: at-least-once in; dedup by event ID BEFORE rollup; warehouse for money
RED FLAGS
- OLAP store as system of record, or source of invoices
- Row-at-a-time inserts into ClickHouse
- Sort key on event_id or time alone in a multi-tenant table
- user_id / session_id as a rollup dimension
- Query-time joins on the user-facing path
- No per-tenant quotas or timeouts
- "Real time" with no lag SLO or data-as-of on the page
- No raw events kept, so no rebuild path