These three are compared because they are the open-source answers to the same request: sub-second aggregations over billions of fresh events, usually straight off Kafka. They overlap on maybe 80% of features — columnar storage, compression, real-time ingestion, SQL — so the feature matrix will not decide it. What decides it is the query population. ClickHouse is a fast, general SQL database you shape with tables and sort keys; it shines when tens to hundreds of concurrent queries ask varied, ad-hoc questions with joins. Druid and Pinot are serving systems you shape with ingestion specs and indexes; they shine when thousands of concurrent users ask the same few questions and each must answer in under 100–300ms. The one question: is this for a few hundred analysts and engineers exploring data, or for thousands of end users loading the same dashboard?
The Verdict#
Default: ClickHouse, because most teams have the internal-analytics problem and ClickHouse is the simplest to run and the most flexible to query. Choose Pinot when the product is user-facing analytics at thousands of QPS with per-tenant queries or upserts; choose Druid when time-series slicing with ingest-time rollup is the core and you already have the expertise.
| Pick ClickHouse when | Pick Druid when | Pick Pinot when |
|---|---|---|
| Users are engineers and analysts: tens to hundreds of concurrent queries | The workload is time-series slicing and dicing with ingest-time rollup | The product is "every user sees their own analytics": 1K–100K QPS |
| Queries are ad hoc and need joins, window functions, full SQL | You want hot data on local disk and history in deep storage, managed by the system | p99 latency must be bounded (star-tree) even for heavy tenants |
| Logs, traces, product analytics, observability backends | Druid expertise exists in-house, or a vendor runs it | You need exact upserts from a stream (order status, latest state per key) |
| One team wants the fewest moving parts: one binary plus a keeper | Rollup can cut rows 10–100× and drill-down to raw events is not required | Many index types per column are needed to keep per-query work tiny |
🎯 Staff Move: "Before I pick an engine I want the query population: how many concurrent queries, how predictable, what p99. For an internal observability store with a few hundred engineers I'd take ClickHouse. For a merchant dashboard serving two million merchants at 10K QPS with a 200ms budget, I'd take Pinot with a tenant-first sort and a star-tree."
When the non-default wins:
- ClickHouse for user-facing analytics wins when concurrency is in the low hundreds of QPS, queries hit materialized views, and a result cache absorbs repeats — one engine for internal and external use is cheaper to run than two.
- Druid over Pinot wins when the team already runs it, and when ingest-time rollup plus compaction of time-series data are the main cost levers rather than per-column index variety or upserts.
- A cloud warehouse instead of all three wins when the users are analysts, concurrency is low, and seconds of latency are fine. Do not run a real-time OLAP cluster for a weekly report.
At a Glance#
| Dimension | ClickHouse | Apache Druid | Apache Pinot |
|---|---|---|---|
| Data model | Tables with a sort key (ORDER BY); MergeTree family engines | Datasources of time-partitioned immutable segments; dimensions and metrics | Tables of segments; realtime, offline and hybrid tables |
| Query language | Full SQL dialect: joins, window functions, CTEs, many functions | SQL (translated to native queries); joins limited; multi-stage engine for heavier SQL | SQL; multi-stage engine adds distributed joins; single-stage is the fast path |
| Indexing | Sparse primary index (one mark per ~8,192-row granule), skip indexes, projections | Bitmap (roaring) indexes on dimensions by default; time partitioning | Sorted, inverted, range, bloom, text, JSON, geo and star-tree indexes per column |
| Pre-aggregation | Materialized views (insert-time triggers), projections, AggregatingMergeTree | Rollup at ingest; compaction can re-roll to coarser grain | Star-tree index with bounded per-query work; rollup on ingest |
| Real-time ingestion | Kafka engine or external consumers; batched inserts (thousands to 100K rows each) | Supervisor-managed Kafka/Kinesis tasks; query-able within seconds | Servers consume Kafka partitions directly; query-able within seconds |
| Upserts / updates | Eventually: ReplacingMergeTree dedups on merge; lightweight deletes and updates | Re-ingest a time interval to replace it | Exact upserts and dedup per primary key, with partitioned input |
| Consistency | Per-replica eventual; inserts atomic per block | Segment versions replace atomically per interval | Per-partition ordering from the stream; upsert masks older docs |
| Ingest throughput | ~100K–1M+ rows/s per node with large batches | ~10K–100K+ rows/s per ingestion task | ~10K–100K+ rows/s per consuming partition |
| Query latency | 10ms–seconds; ad-hoc scans scale with bytes read | 10–500ms for slice-and-dice on rolled-up data | 10–100ms p99 at high QPS with the right indexes |
| Concurrency | Hundreds of QPS per cluster typical; scale with replicas and caching | Thousands of QPS with broker caching | Thousands to tens of thousands of QPS; built for it |
| Scaling model | Shards × replicas; resharding is manual; cloud version separates storage on object store | Separate process types scale independently; deep storage holds everything | Controller, broker, server, minion scale independently; tenants tag servers |
| Operational burden | Lowest: one server binary + ClickHouse Keeper (or ZooKeeper) | Highest: coordinator, overlord, broker, router, historical, middle manager/indexer + metadata DB + deep storage + coordination | High: controller, broker, server, minion + ZooKeeper + deep storage |
| Managed options | ClickHouse Cloud and several third-party hosts | Imply and other hosts | StarTree and other hosts |
| Cost shape | Compute + local or object storage; cheap per TB thanks to 5–20× compression | Historicals hold hot segments locally; deep storage cheap; many node types | Servers sized for index memory and upsert maps; deep storage cheap |
Throughput and latency numbers are orders of magnitude on modern hardware with sensible schemas; any of the three can be 10× slower on a schema that ignores its sort key or partitioning.
Numbers to bring:
| Figure | Value | Condition |
|---|---|---|
| ClickHouse insert batch | 10K–100K rows per insert, or async inserts | Small inserts create too many parts |
| ClickHouse primary index granularity | One mark per ~8,192 rows | Default setting |
| Druid segment size | 300–700 MB, ~5M rows | Documented starting target |
| Rollup reduction | 10–100× for bounded dimensions | Collapses when a high-cardinality field is added |
| Pinot upsert key map | ~(key bytes + 24) bytes per live key | 500M keys ≈ 16 GB before replication |
| Compression | 5–20× vs raw JSON | Typical columnar event data |
| Ingest freshness | 1–30 seconds | Kafka-direct ingestion in all three |
| User-facing p99 target | 100–300ms | Requires pruning by tenant and time |
How They Actually Differ#
1. A Database vs a Serving System#
ClickHouse is a database: you CREATE TABLE ... ENGINE = MergeTree ORDER BY (tenant_id, event_time), insert with SQL, and query with SQL. Its genius is raw scan speed — vectorized execution over compressed columns, often reading billions of rows per second per node — plus a sparse primary index that skips granules outside the sort-key range. You design for it by choosing the sort key and partitioning, then adding materialized views for the expensive aggregates.
Druid and Pinot are serving systems: you write an ingestion spec or table config that declares the time column, dimensions, metrics, rollup or indexes, and the system builds immutable segments, distributes them across servers and routes queries through brokers that scatter and gather. Their genius is predictable latency under concurrency: indexes, segment pruning and pre-aggregation keep each query's work bounded, so 5,000 QPS of tenant-scoped queries does not melt the cluster.
🎯 Staff Insight: ClickHouse's answer to concurrency is "scan faster and cache". Druid's and Pinot's answer is "do less work per query". At 50 QPS the first wins on simplicity; at 10,000 QPS of near-identical queries the second wins on cost.
2. How Each Pays for Pre-Aggregation#
| Technique | ClickHouse | Druid | Pinot |
|---|---|---|---|
| Where it happens | Materialized view triggered on each insert block into a target table | Rollup during ingestion: rows with the same truncated time and dimensions merge | Star-tree built per segment; optional rollup on ingest |
| Effect | Arbitrary aggregates; you query the MV table explicitly | 10–100× fewer rows for bounded dimensions | Hard ceiling on rows touched per query for the chosen dimensions |
| What you lose | Must maintain MVs and backfill them manually | Raw events gone from the store; streaming rollup is best-effort until compaction | Storage and build time; not usable with upsert tables |
| Who pays | The team that owns the schema | Product: no drill-down below the rollup grain | Infra: bigger segments and longer builds |
The rollup ratio is the cost model. 1 billion events/day with (minute, country, device, campaign) at 200 × 5 × 10,000 combinations may collapse 20× — or barely at all if a dimension like user_id sneaks in. Measure it; do not assume it.
3. Freshness, Updates and Correctness#
All three consume Kafka at least once, so duplicates are possible on retries. Their correction stories differ sharply:
- ClickHouse deduplicates identical insert blocks on retry and offers ReplacingMergeTree to keep the latest row per key — but only when parts merge, at an unspecified time. Queries that must be exact use
FINAL(costly) or aggregate withargMax. Lightweight deletes and updates exist; heavy mutation workloads remain an anti-pattern. - Druid treats data as append-only segments per time chunk. To fix a day, you re-ingest that day; the new segment version atomically replaces the old. Corrections take minutes to hours and are batch jobs.
- Pinot has the only exact streaming upsert of the three: an in-memory map from primary key to the latest record, older versions masked at query time. The price is that the input topic must be partitioned by the key, partitions are fixed at table creation, and heap grows with live keys — roughly
keys × (key bytes + 24), so 500M live keys is ~16 GB before replication.
4. Operating Model#
ClickHouse's operational surface is one binary per node plus ClickHouse Keeper (Raft-based and ZooKeeper-protocol compatible) for replication metadata. The pain points are known: too many small inserts creating "too many parts", manual resharding when a cluster grows, and heavy queries from one user starving others without quotas.
Druid has the largest surface: separate coordinator, overlord, broker, router, historical and ingestion processes, plus a metadata database, deep storage and a coordination service. That separation is a strength at large scale — historicals, brokers and ingestion scale independently, and deep storage makes historicals replaceable — and a tax at small scale. Pinot sits between: controller, broker, server and minion processes, with ZooKeeper via Apache Helix for cluster state, and deep storage for segment backup.
| Team situation | Easiest to run |
|---|---|
| 1–2 engineers, one cluster, internal users | ClickHouse |
| Platform team, user-facing product, strict SLOs | Pinot or Druid, or a managed service for either |
| Already running Kafka, ZooKeeper, S3 and Kubernetes well | Any of the three; choose on query population |
Where Each One Breaks#
| System | Failure mode | Symptom | Detection | Mitigation | Owner |
|---|---|---|---|---|---|
| ClickHouse | Too many parts | Inserts rejected; merges fall behind | Active parts per partition, merge queue length | Batch inserts (10K–100K rows), async inserts, coarser partitions | Ingest owner |
| ClickHouse | One query starves the cluster | An ad-hoc full scan pins all cores; dashboards time out | Query duration, memory usage per user | Per-user quotas, max_execution_time, separate replicas for ad-hoc vs serving | Platform |
| ClickHouse | Resharding pain | One shard full or hot; adding shards does not move old data | Disk per shard, rows per shard | Plan sharding key and capacity; or object-storage-backed cloud tier | Platform |
| ClickHouse | Duplicate rows visible | Counts briefly high until merges run | Reconciliation against source | FINAL/argMax for exact queries; idempotent upstream | Data team |
| Druid | Segment explosion | Tens of thousands of tiny segments; query time spent scheduling | Segment count and average size per datasource | Compaction to 300–700 MB segments; coarser segment granularity | Platform |
| Druid | Coordinator or metadata DB trouble | New segments not loaded; data stale though ingestion succeeds | Segment load queue, unavailable segments | Healthy metadata DB, coordinator HA, alert on load lag | Platform |
| Druid | Rollup ratio collapse | Storage and latency jump after a high-cardinality dimension is added | SUM(count)/COUNT(*) ratio trend | Dimension review process; sketches instead of raw IDs | Schema owner |
| Pinot | Upsert memory growth | Server GC storms or OOM; consumption stops on every upsert table on that server | Heap usage, consuming segment lag | Metadata TTL, partial upserts only for hot window, more partitions at creation | Platform |
| Pinot | Hot tenant | One large tenant's queries spike p99 for everyone | Per-tenant latency, server CPU skew | Tenant tagging onto dedicated servers, quotas, star-tree for that table | Platform + product |
| Pinot | Partition count frozen | Ingest cannot scale past topic partitions on upsert tables | Consumer lag | Choose partition count up front; new table + backfill to change it | Platform |
| All three | Query shape drift | A new dashboard filter bypasses the sort key; scans 1,000× more data | Bytes read per query by dashboard | Query review for user-facing paths; result caching; per-query limits | Product + platform |
The production surprise for each:
- ClickHouse: it is so fast on day one that nobody adds quotas, and the first ad-hoc
SELECT *over a year of logs takes down the dashboards. - Druid: the system is healthy, the data is stale. Ingestion succeeded but segments were never handed off and loaded.
- Pinot: the upsert table that was a pilot with 10M keys has 1B keys a year later, and heap is now the capacity limit.
Incident Sketch: The Ad-Hoc Query That Took Down the Dashboards#
t=0 An engineer runs a GROUP BY over 12 months of raw events, no tenant filter
t=+5s Query reads 40 TB across all shards; every core busy; memory near limit
t=+20s Customer dashboard queries on the same replicas queue; p99 from 150ms to 25s
t=+2min Dashboard API times out; support tickets begin
t=+6min On-call kills the query; queue drains in 90 seconds
Prevention: separate replicas (or a separate cluster) for ad-hoc and serving traffic, per-user quotas on memory and execution time, and a rule that user-facing queries must filter on the leading sort-key column. Owner: platform for isolation, product for query review.
Incident Sketch: Stale Druid Dashboard, Green Ingestion#
t=0 Metadata database fails over; coordinator loses its connection briefly
t=+1min Ingestion tasks keep consuming Kafka and publishing segments to deep storage
t=+2min Coordinator stops assigning new segments to historicals
t=+45min Dashboards show data 45 minutes old; ingestion lag metric reads zero
t=+50min Coordinator restarted; load queue drains; data catches up
Detection that would have caught it: a freshness probe that queries the newest timestamp through the broker every minute, plus segment load-queue length and unavailable segment count. Lesson: measure freshness where users read, not where data is written.
Cost and Operations#
| ClickHouse | Druid | Pinot | |
|---|---|---|---|
| Who runs it | A small data-platform team; often 1–2 engineers for a mid-size cluster | A platform team that can run 6+ process types, a metadata DB and deep storage | A platform team; controller and ZooKeeper care plus server sizing |
| What the bill scales with | Bytes stored after compression (5–20× typical) and CPU for scans | Hot segments on historicals (local SSD), ingestion task slots, broker memory for caching | Server memory for indexes and upsert maps, SSD for segments, ingest partitions |
| Cheapest at | Large raw volumes with modest concurrency (logs, traces, events) | High-concurrency time-series with strong rollup | High-concurrency, per-tenant, low-latency queries |
| Cost cliff | Concurrency: adding replicas just to serve more QPS | Process zoo overhead at small scale | Memory for upserts and star-trees at very high key counts |
Back-of-envelope: 1 TB/day of raw events compressed 10× is 100 GB/day; 90 days hot is 9 TB, which fits 3 replicas of a handful of large NVMe nodes in any of the three. The engines diverge on compute for concurrency: serving 10K QPS from ClickHouse means replicas and caches sized for scans, while Pinot or Druid serve it from indexes and rollups with fewer cores per query.
| Scale | Sensible choice | Rough infrastructure | People |
|---|---|---|---|
| Small (100 GB/day, internal users) | Single ClickHouse cluster, 2–3 nodes, or managed | Hundreds to low thousands of dollars a month | Part of one engineer |
| Medium (1 TB/day, internal + first customer dashboards) | ClickHouse with materialized views and cache; separate serving replicas | Low to mid thousands | 1–2 engineers |
| Large (10 TB/day, 10K+ QPS user-facing) | Pinot or Druid for serving, ClickHouse for internal exploration, raw data in object storage | Tens of thousands and up | A platform team of 3–6 |
Assumptions: compressed columnar storage on SSD, 30–90 days hot, list cloud prices; order of magnitude only.
🧭 Principal Insight: The cost line that grows forever is not storage. It is every new dimension a product manager adds to a user-facing dashboard — it changes rollup ratios, index memory and query cost for all tenants. Make dimension additions a reviewed change with a cost estimate.
Switching Later#
| Migration | Difficulty | What's hard to undo |
|---|---|---|
| ClickHouse ↔ Druid ↔ Pinot (data) | Moderate: replay from Kafka or rebuild from raw files in object storage | Nothing, if you kept the raw events; everything, if you only kept rollups |
| Queries and dashboards | Moderate: SQL dialects and functions differ | Engine-specific functions, MV patterns, native query JSON |
| Rollup to raw | Impossible without the raw source | Rolled-up data cannot be un-aggregated |
| Single cluster to tenant-isolated clusters | Moderate | Routing logic in the API; per-tenant SLOs |
| OLAP store to warehouse for billing | Should be the default from day one | Reconciliation contracts with finance |
Moving the user-facing path to a second engine, in order:
- Keep Kafka as the shared source; the new engine consumes the same topics from the retained offset.
- Backfill history from raw files in object storage into offline or batch-ingested segments.
- Shadow the dashboard API: issue every query to both engines, compare results and latency per tenant.
- Cut over tenant cohorts, largest last, with a per-tenant flag to roll back.
- Keep the original engine for internal exploration; delete the duplicated serving tables.
The one-way doors:
- Discarding raw events after rollup. Keep raw data in object storage (Parquet, 30–90 days or longer) so any engine can be rebuilt.
- Pinot upsert table partitioning. The partition count is fixed at creation; changing it means a new table and a backfill.
- Exposing engine-specific SQL to customers. Once customers write queries against your engine's dialect, you cannot switch engines without breaking them. Expose an API, not the SQL.
How Real Companies Chose#
Cloudflare — ClickHouse Replacing a Postgres Rollup Pipeline#
Cloudflare replaced its 2014-era analytics pipeline — a PostgreSQL rollup database and a 12-node Citus cluster fed by thousands of lines of aggregation code — with ClickHouse, after its DNS team had proven it. The new pipeline handled ~6M HTTP requests per second on average (8M peaks) with 36 ClickHouse nodes and 3× replication, reduced stored data from about 60.5 PiB to 18.5 PiB, and raised API throughput from 15 to about 40 queries per second (Cloudflare Blog).
Staff insight: Note the query rate: tens of QPS over enormous data. That is ClickHouse's sweet spot — huge volume, modest concurrency — and the sentence that shows you know where each engine wins.
Confluent — Druid Chosen over Pinot and ClickHouse#
Confluent runs Druid behind its cloud monitoring dashboards, alerting, stream lineage and metrics API, as well as internal billing and control-plane workflows: over three million events per second ingested and over 250 queries per second, with 7 days on historicals and 2 years in S3. In 2019 it evaluated Druid, Pinot and ClickHouse and chose Druid for its documentation at the time, in-house expertise, and proven handling of high-cardinality time-series metrics (Confluent Blog).
Staff insight: The deciding factors were documentation and expertise, not benchmarks. For engines this close in capability, the team's ability to operate it is a legitimate tiebreaker — say so.
Uber — Pinot for User-Facing Real-Time Analytics#
Uber scaled Apache Pinot from a few use cases to a multi-cluster, all-active deployment powering hundreds of use cases over terabyte-scale data with millisecond latency — including restaurant-manager dashboards for Uber Eats, city operations dashboards, and backend services that need fresh analytics, with per-region QPS growing 30× (Uber Engineering). Separately, Uber moved its logging platform from Elasticsearch to ClickHouse, citing a single ClickHouse node ingesting ~300K logs per second and hardware cost cut by more than half (Uber Engineering).
Staff insight: One company, two engines, chosen by query population: Pinot where thousands of external users hit dashboards, ClickHouse where engineers query logs. That is the whole comparison in one org chart.
Follow-Ups to Expect#
| After You Say... | They Will Ask... | What They're Testing |
|---|---|---|
| "ClickHouse" | "Now 5,000 merchants load their dashboard at 9 a.m. What happens?" | Concurrency limits, MVs, caching, replicas vs a serving engine |
| "Pinot for the dashboard" | "One merchant is 1,000× bigger than the median. Then what?" | Tenant isolation, star-tree, quotas |
| "Rollup to the minute" | "Support asks for the raw events from last Tuesday." | Raw retention in object storage; rollup is irreversible |
| "Exact counts" | "Kafka redelivered a batch. Are your numbers wrong?" | At-least-once ingest, dedup mechanisms, reconciliation |
| "Upserts in Pinot" | "How much memory does that take at 1B keys?" | Upsert key map sizing and TTLs |
| "Druid" | "Ingestion is green but the dashboard is stale. Where do you look?" | Segment handoff, coordinator load queue |
| "Use it for billing" | "Would you invoice customers from it?" | Warehouse as source of truth; OLAP numbers are estimates |
| "Add a dimension" | "Product wants to slice by user ID." | Cardinality, rollup collapse, sketches |
| "Sort by tenant then time" | "What about queries across all tenants?" | Separate internal path, pre-aggregated global tables |
| "Materialized views" | "You need to change the aggregation logic. Then what?" | Backfill from raw data; MV versioning |
| "Raw data in object storage" | "How long would a full rebuild take?" | Backfill throughput, parallel ingestion, cost of replay |
What to Say in the Interview#
"The deciding factor is query population, not features. A few hundred engineers running ad-hoc SQL is ClickHouse; tens of thousands of users loading the same dashboard with a 200-millisecond budget is Pinot or Druid."
"I'll sort by tenant then time, roll up to one-minute grain at ingest, and keep raw events in object storage for 30 days, so any rollup decision is reversible by replay."
"Ingestion is at-least-once in all three, so I'll name the dedup mechanism — Pinot upserts if the dashboard must reflect refunds exactly — and keep billing numbers in the warehouse."
"Every new dimension on a customer-facing dashboard is a permanent cost for every tenant, so I'd make adding one a reviewed change with a rollup-ratio estimate."
Related Guides#
- Real-Time OLAP: ClickHouse, Druid & Pinot — internals, indexes, rollups and failure modes in depth
- Design Ad Click Aggregation — a classic real-time aggregation workload
- Design a Metrics Platform — high-cardinality time series at scale
- Time-Series Databases — when a TSDB fits better than general OLAP
- Kafka — the ingestion source for all three
- Batch and Stream Pipelines — hybrid realtime/offline tables and backfills
- Hot Keys — the large-tenant problem in its general form
- Elasticsearch — when the query is text search, not aggregation