Hiring BarSupport

Real-Time OLAP: ClickHouse, Druid & Pinot

Technology guide53 min read6 diagrams

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#

BehaviorSenior (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 modelOne wide table with every fieldSort key from the dominant filter, rollup granularity from the product's slowest acceptable drill-downWrites 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 SLOPrices freshness: the delta between 5-minute batch and 5-second streaming, per product
CorrectnessCounts what arrivesNames the dedup and upsert mechanism, and that background merges do not guarantee itDecides which numbers are billing-grade (warehouse of record) and which are directional (OLAP), and writes that on the dashboard
CostNode countBytes 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 buyPicks the engine they knowPicks by concurrency, upsert needs, and join needsSelf-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#

IntentConstraintStrategyFailure ModeCorrectness Bar
User-facing analytics (merchant dashboards, creator stats, "who viewed my profile")1K–50K QPS, p99 < 100–300 ms, every query scoped to one tenantPinot/Druid or ClickHouse with materialized views; tenant-first sort key; star-tree or rollups; result cacheOne large tenant's scan saturates shared servers; p99 collapses for everyoneDirectional; 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 alertsStream ingest, short hot retention (7–30 days), approximate distinct countsIngest lag during a spike hides the incident the dashboard exists to showApproximate is fine; staleness must be visible
Internal ad-hoc exploration (product analytics, log analytics)10–100 QPS, unpredictable queries, joins, long retentionClickHouse or a cloud warehouse; raw events, wide tables, projectionsOne analyst's unbounded query burns the cluster; nobody owns the schemaReproducible; 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#

PositionRationale
Design from the query, not the eventList 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 indexAll 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 elsewherePre-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 unitStorage 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 defaultIngest 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 addedA 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 recordIt 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:

EncodingWorks OnTypical Effect
DictionaryStrings with bounded distinct values (country, device)2–4 byte integer IDs instead of strings; enables bitmap indexes
Run-lengthSorted, low-cardinality leading columns (tenant_id when sorted first)A million identical values stored as one run
Delta / double-deltaTimestamps, monotonic IDsA few bits per value
Bit-packingSmall integers after dictionary encoding7 bits for 100 distinct values instead of 32
General codec (LZ4, ZSTD)Everything, on top of the aboveAnother 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_id first, 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.

ConceptClickHouseApache DruidApache Pinot
Unit of dataPart (immutable, merged in background) inside a partitionSegment (time-chunked, typically 300–700 MB)Segment (realtime "consuming" or completed, offline)
Ingest from KafkaKafka table engine + materialized view, or an external consumer doing batched INSERTKafka indexing service: supervisor spawns ingest tasksServers consume Kafka partitions directly; segments seal and upload to deep store
Query routingDistributed table fans out to shardsBroker → historicals + realtime tasksBroker → servers, prunes by time, partition, metadata
DurabilityReplication via ClickHouse Keeper (ZooKeeper-compatible); local disks or object storageDeep storage (S3/HDFS) is the source of truth; historicals cacheDeep store holds completed segments; servers replicate consuming segments
Pre-aggregationMaterialized views into AggregatingMergeTree / SummingMergeTree, projectionsRollup at ingest is a first-class table settingStar-tree index; ingestion aggregation
Mutable dataReplacingMergeTree (eventual), lightweight deletes, mutations (heavy)Re-ingest the time interval (segment replacement)Full and partial upsert on realtime tables; dedup
CoordinationKeeper for replication onlyZooKeeper + metadata DB (Postgres/MySQL)Helix on ZooKeeper
Sweet spotAd-hoc SQL, joins, log and product analytics, one team running itTime-series slice-and-dice with rollup, high-concurrency dashboardsVery 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#

Diagram: The Write Path — Kafka to Queryable in Seconds
  1. 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.
  2. 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.
  3. Serve. Historical servers load sealed segments (memory-mapped). The query layer merges results from consuming and sealed segments.
  4. 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#

Diagram: 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.

IndexEngineHow It HelpsCost
Sparse primary index (sort key)ClickHouseOne 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 columnPinot, Druid (time)Range of row IDs for a value: O(log n)One per table
Bitmap / invertedDruid (default on dimensions), PinotRoaring bitmap per value; AND/OR across filtersGrows with cardinality; useless for unique IDs
RangePinotNumeric range filters without scanningExtra storage
Bloom filter / skip indexAll threeSkip blocks for point-ish lookups (request_id = ...)False positives; tuned per column
Star-treePinotPre-aggregated tree over chosen dimensions; bounded work per queryStorage; build time; not usable with upsert
Projection / materialized viewClickHouseSecond sort order or pre-aggregated copy, chosen automatically or by queryDoubles 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).

Diagram: The Star-Tree Index — Pre-Aggregation You Can Tune

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 WHERE clauses 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 ShapeFiltersGroup ByFreshnessShare of QPS
Q1Sales today vs yesterdaytenant, time (2 days)hour< 1 min40%
Q2Top 10 products, last 7 daystenant, timeproduct_id< 5 min25%
Q3Sales by country / channeltenant, timecountry, channel< 5 min20%
Q4Unique buyers, last 30 daystenant, time— (distinct)< 1 h10%
Q5Order drill-down (raw rows)tenant, order_id or time window—< 1 min5%

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:

  1. 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 in ORDER BY instead; it also says partitioning coarser than a month is rarely needed (MergeTree docs).
  2. 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.
  3. 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.
  4. 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 ReplacingMergeTree collapses it, so a replayed batch is counted twice in the rollup. Dedup must happen upstream of the MV, not in the raw table.
  5. Distinct counts as sketch state. uniqState stores 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#

StrategyWhat WorksWhat BreaksWho Pays
Raw rows, aggregate at queryFull flexibility; any new question answerable; simplest pipelineCost scales with rows scanned × QPS; p99 grows with dataInfra budget; users at peak
Rollup at ingest (Druid rollup, ClickHouse MV)20–100× fewer rows for bounded dimensions; predictable latencyRaw detail gone from the store; dimensions fixed at ingest; streaming rollup is best-effortProduct: no drill-down below the grain; data team for backfills
Index-time pre-aggregation (star-tree, projections)Rollup speed and raw rows keptStorage 1.2–2× more; fixed aggregation list; rebuild on changeStorage budget; ingest CPU
Pre-computed result tables (batch or stream job writes final answers)Lowest latency; trivially cacheableEvery new question is an engineering ticketData 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.

Diagram: Upserts and Deduplication
EngineMechanismGuaranteePrice
Pinot upsert (full or partial)In-memory map from primary key to the latest record location; older versions are masked at query timeLatest version per key, exact, within a partitionInput 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 dedupDrop records whose primary key was already seenFirst write wins, exactSame partitioning and key-map memory
ClickHouse ReplacingMergeTreeKeeps the row with the highest version when parts mergeEventual: 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 dedupReplicated tables drop identical retried insert blocksRetries of the same block, not semantic duplicatesRetry must resend the same batch
DruidRe-ingest the affected time chunk; new segment version atomically replaces the oldExact for the replaced intervalBatch 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:

  1. Destroys rollup. Every row becomes unique; the ratio goes to 1.
  2. Bloats dictionaries and bitmaps. A dictionary with 500M entries per segment and a bitmap per value cost more than the column they index.
  3. Explodes group-by state. GROUP BY url over a day can produce 100M groups per server, all held in memory before the merge.

What to do instead, ranked by preference:

NeedPattern
Count distinct usersHLL / Theta sketch metric, mergeable across segments (1–2% error)
Top-N URLsApproximate top-K, or a separate rollup table keyed by URL with daily grain
Look up one user's eventsRaw table sorted by (tenant_id, user_id, time) with short retention, or the OLTP store
Filter by a specific IDBloom filter / skip index, not bitmap or dictionary-backed inverted index
Group by a long-tail fieldBucket 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.

SettingFreshnessQuery flexibilityCost per queryWhen
Raw stream ingest, no rollup5–30 sFullHighestOps/debug dashboards with tens of QPS
Stream ingest + minute rollup5–30 sFixed dimensionsLowAd pacing, merchant "today" tiles
Stream head + hourly batch tail (hybrid)5–30 s head, hours tailFixed in tailLowest at scaleUser-facing products over months of history
Batch ingest only, every 15–60 min15–60 minFixedLowest; simplest opsReports that nobody reads within the hour
Warehouse, no OLAP store5 min to hoursFull SQL, joinsSeconds and $/TB scannedInternal 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#

DimensionClickHouseApache DruidApache PinotCloud Warehouse (BigQuery, Snowflake, Redshift)
Designed forFast SQL analytics on one table at a timeInteractive slice-and-dice over event streamsUser-facing analytics at very high QPSAd-hoc analysis, ELT, joins across the business
Typical p9910 ms – seconds, depends on modeling50 ms – 1 s10–100 ms for tuned queriesSeconds to minutes
Concurrency comfortHundreds of QPS per cluster before careful tuningThousandsThousands to tens of thousandsTens, with queueing and per-query billing
Stream freshnessSeconds (Kafka engine or batched inserts)Seconds (Kafka indexing service)Seconds (direct partition consumption)Minutes typically; streaming inserts exist at extra cost
JoinsGood: hash joins, dictionaries; large joins need careLimited; lookups and broadcast joinsLimited; lookup joins, multi-stage engineFull
MutabilityEventual via merges; mutations are heavyInterval re-ingestUpsert and dedup on realtime tablesFull DML
Ops footprintSmallest: one binary + KeeperLargest: 5–6 process types + ZooKeeper + metadata DB + deep storageMedium: controller, broker, server, minion + ZooKeeper + deep storeNone (managed)
Cost shapeHardware you runHardware you runHardware 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#

Diagram: 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#

QuantityPlanning 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 → consumer1:1; ~10–50K events/s per consumer is a safe planning range
Druid segment target300–700 MB, ~5M rows
ClickHouse insert batch≥ 1,000 rows, ideally 10K–100K
ClickHouse granule8,192 rows default
Star-tree leaf threshold10,000 records default
Pinot upsert key map~(key bytes + 24) bytes per live key
Freshness, streaming ingest5–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 AxisFirst Thing to BreakMove
Event rateConsumers fall behind; part or segment count explodesMore partitions (planned up front), bigger batches, compaction capacity
Query QPSBroker merge CPU and fan-out to many serversPartition-aware routing so a tenant hits 1–2 servers; replica groups; result cache
Data historySegment count and metadata; cold queries hit object storageTiered storage; coarser rollups for old data; fewer, bigger segments
TenantsOne giant tenant sets everyone's p99Tenant tiering; quotas; dedicated tables
DimensionsRollup ratio collapses; dictionaries bloatDimension 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_seconds per partition; ingest.max_event_time_age per 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_count per 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_seconds and bytes_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; ReplacingMergeTree on 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_count per 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#

FailureDetection SignalBlast RadiusMitigationOwner
Ingest lag spiralconsumer lag seconds; event-time ageAll dashboards on the tableScale consumers, shed enrichments, bannerData platform
Part/segment explosionpart and segment countsWhole table, then clusterBatch, compact, pause writerData platform
Noisy tenantper-tenant CPU and bytes scannedEvery tenant on shared serversKill, quota, tierQuery service team
Double countingreconciliation driftBusiness decisions, invoices if misusedRe-aggregate the intervalData platform + finance
Upsert OOMheap, key countEvery upsert table on the replica groupTTL, scale, quarantineData platform
Deep storage outageupload failuresNew segments stay on servers; disks fillExtend local retention; pause compactionData platform + storage

When to Use vs. Alternatives#

SituationPickWhy
End-user analytics, thousands of QPS, fixed queries, seconds freshnessPinot 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 QPSClickHouseSQL flexibility, joins, smallest ops footprint
Analysts, cross-domain joins, minutes freshness fineThe warehouse you already pay forNo new system to run
Infra metrics with labels and alertingTSDB (Time-Series Databases)PromQL, alerting, rate/histogram semantics
Full-text search over logsElasticsearchInverted index on text, relevance
Under ~100M–1B rows, low QPS, one teamPostgreSQL with rollup tablesFree in operational terms; you already run it
Money, invoices, payoutsWarehouse of record or OLTP with reconciliationExactness 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 123 in 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#

MetricAlert Threshold (example)Why
ingest.event_time_age_seconds> 120 s for 5 minThe freshness promise
kafka_consumer_lag_seconds> 60 s and risingLeading indicator of the above
query.latency_p99_ms per query template> 2× baseline for 10 minUser-facing SLO
query.bytes_scanned per tenanttop tenant > 30% of clusterNoisy-neighbour early warning
parts_per_partition / segment_countClickHouse > 300 active parts in a partitionPrecursor to insert failures
rollup_ratiodrops > 30% week over weekA new dimension or a cardinality bug
reconciliation_drift_pct> 0.5% hourlyDuplicates or loss
upsert_primary_keys per server> 80% of budgetOOM 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 StudyHow Real-Time OLAP Is UsedKey Pattern
Ad Click AggregationDruid or Pinot serving per-campaign counts to advertiser dashboards and pacingDedup upstream, rollup at ingest, warehouse reconciles billing
Stream ProcessingServing store downstream of Flink windows, upserted by (campaign, window_start)Upsert key = window key; replay overwrites, never adds
Metrics & Alerting PlatformHigh-cardinality questions (per customer, per endpoint) that a TSDB refusesTSDB for alerting; OLAP for exploration
Leaderboard"Top products this quarter" style analytics boardsRedis for live ranks; OLAP for ad-hoc slicing
Database SelectionThe derived analytical store fed by CDCNever run scans on the OLTP primary
URL ShortenerPer-link click analytics behind the redirect pathAsync event pipeline; analytics never on the hot path
Search IndexingWhere aggregation-heavy queries go instead of the search clusterSearch 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 SayWhat Interviewers HearWhat 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 analystsFive process types for a warehouse job"That's the warehouse. Real-time OLAP when concurrency demands it."

L5 vs L6 vs L7 Responses#

ScenarioL5 AnswerL6 / Staff AnswerL7 / Principal Answer
"Build merchant analytics for 2M sellers"Events into ClickHouse, dashboard queries itQuery shapes first; tenant-first sort key; minute rollup; HLL for uniques; query service with templates, quotas, cache; lag bannerOne 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 serversFind the tenant and the query by bytes scanned; quota it; tier the top tenants; fix the templateIntroduces per-tenant cost budgets and makes the largest customers' tier a priced product
"Numbers disagree with finance"Re-run the queryReconciliation job, find the replay, move dedup before the rollupPublishes which tables are certified, the tolerance finance accepted, and who signs off on corrections
"Should we adopt Pinot?"Yes, it's what LinkedIn usesOnly if QPS > ~1K or upserts needed; otherwise ClickHouse or the warehouseModels 3-year cost including the platform team, names the exit path, and sets the trigger to revisit

The Staff Real-Time OLAP Checklist#

  1. Query population: "Who queries, how many at once, at what p99? Analysts at 20 QPS or users at 5,000?"
  2. Query shapes: "Five templates: filters, group-bys, time ranges. The sort key is (tenant_id, time)."
  3. Grain and retention: "Raw 14 days, minute rollup 13 months; product signed off on no single-order drill-down after two weeks."
  4. Correctness: "At-least-once ingest, dedup by event ID upstream of the rollup, sketches for distinct counts, warehouse for anything billed."
  5. Freshness: "Kafka-direct ingest, 30-second target, event-time age on the dashboard, page at 2 minutes."
  6. 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#

DimensionSenior (L5)Staff (L6)Principal (L7)
ModelingWide table of eventsSort key, grain, rollup and indexes derived from named query shapesDimension admission policy; schema changes reviewed with a backfill budget
CorrectnessAssumes counts are rightPlaces dedup before aggregation; knows merge-time dedup is eventual; uses sketches knowinglyCertified 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 tenantsRebuild-from-source drill; regional DR via re-ingest; platform staffing
ChoiceThe engine they knowPicks by concurrency, upserts, joins, ops footprintBuild vs buy, exit path, when to keep the warehouse

Strong hire signals

SignalWhat 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

SignalWhy It Misses the Bar
No sort key or grain discussedWill scan everything and buy hardware to compensate
Exact distinct counts everywhereUnbounded memory and fan-out cost
OLAP as the source of truthNo rebuild path; correctness unowned
No tenant isolation in a multi-tenant productOne 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.

OptionWhat WorksWhat BreaksWho Pays
One shared platform (one engine, multi-tenant clusters)One expert team; shared ingest and tooling; consistent metric definitionsOne team's dimension or query hurts everyone; platform becomes a bottleneck for schema changesPlatform team and every product during noisy-neighbour incidents
Shared platform with tiers and chargebackIsolation for heavy users; cost visible to the teams creating itNeeds metering, quotas and an admission processPlatform headcount for the control plane
Store per product teamAutonomy; fast iteration3 engines, 6 clusters, inconsistent definitions of "active user", nobody to page at 3amFinance through duplication; data consumers through conflicting numbers
Warehouse only, no real-time storeLowest ops; single sourceCannot serve user-facing QPS or seconds-fresh dashboardsProduct: 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_id to 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.

ScaleWorkloadServers + storage/monthKafka + pipelines/monthPlatform headcountTotal/month
Startup50M events/day, 200 QPS, one product3 nodes ~$6K~$2K0.5 FTE (~$10K)~$18K, versus ~$3–8K for warehouse plus cached queries
Growth5B events/day, 3K QPS, 5 products16 nodes + 100 TB object ~$35K~$10K2 FTE (~$42K)~$87K
Enterprise100B events/day, 20K QPS, 30 products120 nodes + 2 PB object ~$290K~$70K7 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#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionReversibilityCost to Reverse
Kafka partition count for upsert tablesOne-way for that tableNew topic, new table, full re-ingest
Rollup grain with no raw events keptOne-wayLost detail is gone; only future data can be finer
Sort key / partition key / star-tree split orderOne-way-ishRebuild from raw events: days of compute, an engineer-week
Exposing OLAP numbers in a customer-facing API contractOne-way-ishCustomers build on its semantics and freshness
Engine choice (ClickHouse vs Pinot vs Druid)One-way-ishQuarters: ingestion specs, query templates, operations runbooks
Retention per tierTwo-way going shorter, one-way for data already deletedStorage cost only, until something is deleted
Adding a nullable column, a new MV, a new indexTwo-wayBackfill time
Result cache TTL, query quotasTwo-wayConfig 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#

SignalWhat 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
  1. Loading the index…