Hiring BarSupport

ClickHouse vs Druid vs Pinot

Comparison19 min read4 diagrams

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 whenPick Druid whenPick Pinot when
Users are engineers and analysts: tens to hundreds of concurrent queriesThe workload is time-series slicing and dicing with ingest-time rollupThe product is "every user sees their own analytics": 1K–100K QPS
Queries are ad hoc and need joins, window functions, full SQLYou want hot data on local disk and history in deep storage, managed by the systemp99 latency must be bounded (star-tree) even for heavy tenants
Logs, traces, product analytics, observability backendsDruid expertise exists in-house, or a vendor runs itYou need exact upserts from a stream (order status, latest state per key)
One team wants the fewest moving parts: one binary plus a keeperRollup can cut rows 10–100× and drill-down to raw events is not requiredMany 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."

Diagram: The Verdict

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#

DimensionClickHouseApache DruidApache Pinot
Data modelTables with a sort key (ORDER BY); MergeTree family enginesDatasources of time-partitioned immutable segments; dimensions and metricsTables of segments; realtime, offline and hybrid tables
Query languageFull SQL dialect: joins, window functions, CTEs, many functionsSQL (translated to native queries); joins limited; multi-stage engine for heavier SQLSQL; multi-stage engine adds distributed joins; single-stage is the fast path
IndexingSparse primary index (one mark per ~8,192-row granule), skip indexes, projectionsBitmap (roaring) indexes on dimensions by default; time partitioningSorted, inverted, range, bloom, text, JSON, geo and star-tree indexes per column
Pre-aggregationMaterialized views (insert-time triggers), projections, AggregatingMergeTreeRollup at ingest; compaction can re-roll to coarser grainStar-tree index with bounded per-query work; rollup on ingest
Real-time ingestionKafka engine or external consumers; batched inserts (thousands to 100K rows each)Supervisor-managed Kafka/Kinesis tasks; query-able within secondsServers consume Kafka partitions directly; query-able within seconds
Upserts / updatesEventually: ReplacingMergeTree dedups on merge; lightweight deletes and updatesRe-ingest a time interval to replace itExact upserts and dedup per primary key, with partitioned input
ConsistencyPer-replica eventual; inserts atomic per blockSegment versions replace atomically per intervalPer-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 latency10ms–seconds; ad-hoc scans scale with bytes read10–500ms for slice-and-dice on rolled-up data10–100ms p99 at high QPS with the right indexes
ConcurrencyHundreds of QPS per cluster typical; scale with replicas and cachingThousands of QPS with broker cachingThousands to tens of thousands of QPS; built for it
Scaling modelShards × replicas; resharding is manual; cloud version separates storage on object storeSeparate process types scale independently; deep storage holds everythingController, broker, server, minion scale independently; tenants tag servers
Operational burdenLowest: one server binary + ClickHouse Keeper (or ZooKeeper)Highest: coordinator, overlord, broker, router, historical, middle manager/indexer + metadata DB + deep storage + coordinationHigh: controller, broker, server, minion + ZooKeeper + deep storage
Managed optionsClickHouse Cloud and several third-party hostsImply and other hostsStarTree and other hosts
Cost shapeCompute + local or object storage; cheap per TB thanks to 5–20× compressionHistoricals hold hot segments locally; deep storage cheap; many node typesServers 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:

FigureValueCondition
ClickHouse insert batch10K–100K rows per insert, or async insertsSmall inserts create too many parts
ClickHouse primary index granularityOne mark per ~8,192 rowsDefault setting
Druid segment size300–700 MB, ~5M rowsDocumented starting target
Rollup reduction10–100× for bounded dimensionsCollapses when a high-cardinality field is added
Pinot upsert key map~(key bytes + 24) bytes per live key500M keys ≈ 16 GB before replication
Compression5–20× vs raw JSONTypical columnar event data
Ingest freshness1–30 secondsKafka-direct ingestion in all three
User-facing p99 target100–300msRequires 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.

Diagram: 1. A Database vs a Serving System

🎯 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#

TechniqueClickHouseDruidPinot
Where it happensMaterialized view triggered on each insert block into a target tableRollup during ingestion: rows with the same truncated time and dimensions mergeStar-tree built per segment; optional rollup on ingest
EffectArbitrary aggregates; you query the MV table explicitly10–100× fewer rows for bounded dimensionsHard ceiling on rows touched per query for the chosen dimensions
What you loseMust maintain MVs and backfill them manuallyRaw events gone from the store; streaming rollup is best-effort until compactionStorage and build time; not usable with upsert tables
Who paysThe team that owns the schemaProduct: no drill-down below the rollup grainInfra: 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 with argMax. 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.
Diagram: 3. Freshness, Updates and Correctness

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 situationEasiest to run
1–2 engineers, one cluster, internal usersClickHouse
Platform team, user-facing product, strict SLOsPinot or Druid, or a managed service for either
Already running Kafka, ZooKeeper, S3 and Kubernetes wellAny of the three; choose on query population

Where Each One Breaks#

SystemFailure modeSymptomDetectionMitigationOwner
ClickHouseToo many partsInserts rejected; merges fall behindActive parts per partition, merge queue lengthBatch inserts (10K–100K rows), async inserts, coarser partitionsIngest owner
ClickHouseOne query starves the clusterAn ad-hoc full scan pins all cores; dashboards time outQuery duration, memory usage per userPer-user quotas, max_execution_time, separate replicas for ad-hoc vs servingPlatform
ClickHouseResharding painOne shard full or hot; adding shards does not move old dataDisk per shard, rows per shardPlan sharding key and capacity; or object-storage-backed cloud tierPlatform
ClickHouseDuplicate rows visibleCounts briefly high until merges runReconciliation against sourceFINAL/argMax for exact queries; idempotent upstreamData team
DruidSegment explosionTens of thousands of tiny segments; query time spent schedulingSegment count and average size per datasourceCompaction to 300–700 MB segments; coarser segment granularityPlatform
DruidCoordinator or metadata DB troubleNew segments not loaded; data stale though ingestion succeedsSegment load queue, unavailable segmentsHealthy metadata DB, coordinator HA, alert on load lagPlatform
DruidRollup ratio collapseStorage and latency jump after a high-cardinality dimension is addedSUM(count)/COUNT(*) ratio trendDimension review process; sketches instead of raw IDsSchema owner
PinotUpsert memory growthServer GC storms or OOM; consumption stops on every upsert table on that serverHeap usage, consuming segment lagMetadata TTL, partial upserts only for hot window, more partitions at creationPlatform
PinotHot tenantOne large tenant's queries spike p99 for everyonePer-tenant latency, server CPU skewTenant tagging onto dedicated servers, quotas, star-tree for that tablePlatform + product
PinotPartition count frozenIngest cannot scale past topic partitions on upsert tablesConsumer lagChoose partition count up front; new table + backfill to change itPlatform
All threeQuery shape driftA new dashboard filter bypasses the sort key; scans 1,000× more dataBytes read per query by dashboardQuery review for user-facing paths; result caching; per-query limitsProduct + 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#

ClickHouseDruidPinot
Who runs itA small data-platform team; often 1–2 engineers for a mid-size clusterA platform team that can run 6+ process types, a metadata DB and deep storageA platform team; controller and ZooKeeper care plus server sizing
What the bill scales withBytes stored after compression (5–20× typical) and CPU for scansHot segments on historicals (local SSD), ingestion task slots, broker memory for cachingServer memory for indexes and upsert maps, SSD for segments, ingest partitions
Cheapest atLarge raw volumes with modest concurrency (logs, traces, events)High-concurrency time-series with strong rollupHigh-concurrency, per-tenant, low-latency queries
Cost cliffConcurrency: adding replicas just to serve more QPSProcess zoo overhead at small scaleMemory 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.

ScaleSensible choiceRough infrastructurePeople
Small (100 GB/day, internal users)Single ClickHouse cluster, 2–3 nodes, or managedHundreds to low thousands of dollars a monthPart of one engineer
Medium (1 TB/day, internal + first customer dashboards)ClickHouse with materialized views and cache; separate serving replicasLow to mid thousands1–2 engineers
Large (10 TB/day, 10K+ QPS user-facing)Pinot or Druid for serving, ClickHouse for internal exploration, raw data in object storageTens of thousands and upA 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#

MigrationDifficultyWhat's hard to undo
ClickHouse ↔ Druid ↔ Pinot (data)Moderate: replay from Kafka or rebuild from raw files in object storageNothing, if you kept the raw events; everything, if you only kept rollups
Queries and dashboardsModerate: SQL dialects and functions differEngine-specific functions, MV patterns, native query JSON
Rollup to rawImpossible without the raw sourceRolled-up data cannot be un-aggregated
Single cluster to tenant-isolated clustersModerateRouting logic in the API; per-tenant SLOs
OLAP store to warehouse for billingShould be the default from day oneReconciliation contracts with finance
Diagram: Switching Later

Moving the user-facing path to a second engine, in order:

  1. Keep Kafka as the shared source; the new engine consumes the same topics from the retained offset.
  2. Backfill history from raw files in object storage into offline or batch-ingested segments.
  3. Shadow the dashboard API: issue every query to both engines, compare results and latency per tenant.
  4. Cut over tenant cohorts, largest last, with a per-tenant flag to roll back.
  5. Keep the original engine for internal exploration; delete the duplicated serving tables.

The one-way doors:

  1. Discarding raw events after rollup. Keep raw data in object storage (Parquet, 30–90 days or longer) so any engine can be rebuilt.
  2. Pinot upsert table partitioning. The partition count is fixed at creation; changing it means a new table and a backfill.
  3. 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."

  1. Loading the index…