Apache Druid is a database for one kind of question asked very often: how many, how much, broken down by what, over which time range, right now. Product analytics dashboards, network telemetry, ad-tech reporting and operational monitoring all ask that question thousands of times a minute against billions of events that are still arriving. Druid answers it in well under a second by organising everything around time, storing data in immutable columnar segments with indexes on every string column, and spreading those segments across many servers that scan them in parallel.

That design is also why Druid surprises people who treat it like a general database. Rows are not updated in place, joins are limited compared with a warehouse, and the cluster is a collection of specialised services with three external dependencies. This article explains each part, follows an event from a Kafka topic to a query result, works through a sizing example, and lists the failure modes operators actually meet.

Advertisement

The services and the dependencies

Druid splits its work into services that you can run together on a few machines or separately at scale. The documentation groups them into three server types.

  • Master servers run the Coordinator, which decides which data servers hold which segments and keeps them balanced, and the Overlord, which accepts ingestion work and assigns tasks.
  • Query servers run the Broker, which receives queries, works out which segments are involved and where they live, fans the query out and merges the partial results, and optionally the Router, a front door that routes requests to Brokers, Coordinators and Overlords and serves the web console.
  • Data servers run Historicals, which store and serve published segments, and the ingestion workers: either a MiddleManager that forks a separate Peon process per task, or the newer Indexer, which runs tasks as threads in one JVM. You deploy one style or the other, not both.

Three external systems hold the state. Deep storage, usually S3, HDFS or another object store, holds the durable copy of every segment. The metadata store, usually PostgreSQL or MySQL, records which segments exist and are in use, plus task state, supervisor specs and rules. ZooKeeper is used for service discovery, coordination and leader election among the masters; the current documentation still lists it as a required dependency, and the role ZooKeeper plays in Hadoop clusters is a good primer.

The important consequence of this split is that data servers are disposable. A Historical holds only a cache of segments whose real home is deep storage, so losing one means the Coordinator asks other Historicals to load its segments again. Losing deep storage or the metadata store is the serious event, which is why those two deserve the same care as any production database.

Druid: master, query and data servers around three external dependenciesClientsSQL over HTTP, JDBCRouteroptional front doorBrokerplan, scatter, mergeCoordinatorsegment placementOverlordingestion tasksHistoricalsserve published segmentsMiddleManager / Peonsor Indexer: ingest, serve freshqueryload / dropassign tasksKafka / Kinesisevent streamsconsumeDeep storageS3, HDFS, GCSMetadata storeMySQL, PostgreSQLZooKeeperdiscovery, leaderspublishpullsegment recordsHistoricals copy segments from deep storage to local disk and memory-map them; deep storage is the durable copyFresh data is queryable on the ingestion tasks until handoff moves it to Historicals
Brokers plan and merge queries, Historicals and ingestion tasks scan segments, the masters place data and run tasks, and deep storage plus the metadata store hold the durable state.

The segment: where the speed comes from

Every Druid table, called a datasource, is partitioned first by time into chunks set by the segment granularity, such as one hour or one day. Each time chunk holds one or more segment files, typically a few hundred megabytes and a few million rows each. A segment is identified by its datasource, interval, version and partition number, and it is immutable once published. To change data you write new segments for the same interval with a newer version; when they are published, queries switch to them atomically and the older ones become unused. This versioning is how Druid offers reindexing, compaction and batch replacement without locking readers.

Inside a segment, data is stored by column. The __time column holds the timestamp. String dimensions are dictionary encoded: each distinct value gets an integer id, the column stores ids, and for each value Druid keeps a bitmap of the rows that contain it. A filter such as country = 'DE' AND device = 'mobile' becomes an AND of two compressed bitmaps, often without touching the rows at all; bitmap indexes explains why that is so cheap. Numeric metrics are stored as compressed arrays. Reading only the referenced columns is the usual columnar advantage.

The second source of speed is rollup. If you enable it and set a query granularity of, say, one minute, Druid aggregates rows at ingestion time that share the same truncated timestamp and the same values of every dimension, storing counts and sums instead of raw events. A clickstream with 500 raw events per minute for one page, country and device becomes one row. Rollup is a trade: you can no longer see individual events or filter on columns you did not keep, and sketches must be used for distinct counts. Its effectiveness depends on dimension cardinality, so a high-cardinality column such as a user id or request id can destroy it entirely.

Advertisement

Streaming ingestion and handoff

For Kafka or Kinesis you submit a supervisor spec to the Overlord. The supervisor runs continuously, launches indexing tasks, assigns them partitions and replaces them when they finish or fail.

{
  "type": "kafka",
  "spec": {
    "dataSchema": {
      "dataSource": "clicks",
      "timestampSpec": {"column": "ts", "format": "iso"},
      "dimensionsSpec": {"dimensions": ["page", "country", "device", "referrer"]},
      "metricsSpec": [
        {"type": "count", "name": "events"},
        {"type": "longSum", "name": "bytes", "fieldName": "bytes"}
      ],
      "granularitySpec": {"segmentGranularity": "hour", "queryGranularity": "minute", "rollup": true}
    },
    "ioConfig": {
      "topic": "clicks",
      "inputFormat": {"type": "json"},
      "consumerProperties": {"bootstrap.servers": "kafka-1:9092"},
      "taskCount": 4,
      "replicas": 2,
      "taskDuration": "PT1H"
    },
    "tuningConfig": {"type": "kafka"}
  }
}

Each task reads its partitions, builds rows in memory, persists them to local disk periodically, and answers queries for the data it holds, which is what makes events queryable within seconds of arriving. When the task duration ends, the task merges its data into segments, pushes them to deep storage and publishes them by writing segment records to the metadata store. It commits the Kafka offsets in the same metadata transaction as the segment records, which is how Druid gets exactly-once ingestion from a stream: segments and offsets either both commit or neither does.

Then comes handoff. The Coordinator notices the new segments, assigns them to Historicals by the datasource's load rules, and the Historicals download, memory-map and announce them. Only when a Historical announces a segment does the task stop serving it and exit. Setting replicas to 2 runs two tasks reading the same partitions, so a failed task does not create a gap in fresh data. For Kafka basics see Kafka in streaming systems.

For batch loads, the SQL-based multi-stage query engine, introduced in Druid 24, is the usual path. It reads external files and writes segments in a single statement, and the same engine can also query segments directly from deep storage.

REPLACE INTO clicks OVERWRITE WHERE __time >= TIMESTAMP '2026-09-30' AND __time < TIMESTAMP '2026-10-01'
SELECT TIME_PARSE(ts) AS __time, page, country, device, referrer,
       COUNT(*) AS events, SUM(bytes) AS bytes
FROM TABLE(EXTERN(
  '{"type":"s3","prefixes":["s3://raw/clicks/2026-09-30/"]}',
  '{"type":"json"}',
  '[{"name":"ts","type":"string"},{"name":"page","type":"string"},{"name":"country","type":"string"},
    {"name":"device","type":"string"},{"name":"referrer","type":"string"},{"name":"bytes","type":"long"}]'))
GROUP BY 1, 2, 3, 4, 5
PARTITIONED BY DAY
CLUSTERED BY page

The query path

A query arrives at a Broker as SQL or as native JSON. The Broker plans SQL into native queries such as timeseries, topN, groupBy or scan. It then consults its view of the cluster, built from announcements, which maps every segment to the servers holding it. The time filter prunes that list first: a query over the last hour touches only one or two hourly time chunks no matter how much history the datasource holds. Secondary partitioning, such as range partitioning on a dimension, can prune further.

The Broker sends each selected segment's work to one replica, on a Historical for published data or on an ingestion task for fresh data, and those servers scan their segments in parallel using bitmaps for filters and only the needed columns. Each returns partial aggregates, and the Broker merges them, applies final ordering and limits, and returns the result. Per-segment results from Historicals can be cached, which helps dashboards that repeat the same query, because immutable segments never invalidate.

SELECT TIME_FLOOR(__time, 'PT5M') AS bucket, country, SUM(events) AS events
FROM clicks
WHERE __time >= CURRENT_TIMESTAMP - INTERVAL '1' HOUR AND device = 'mobile'
GROUP BY 1, 2
ORDER BY events DESC
LIMIT 50

The interactive engine is built for many fast, mostly aggregating queries. Large joins, high-cardinality groupings and deep subqueries are expensive because results are merged on the Broker. The multi-stage engine handles heavier work, and Druid 31 added an experimental engine called Dart aimed at exactly those complex queries. Treat it as experimental until your version's documentation says otherwise.

Worked example: sizing a clickstream datasource

Suppose a site produces 2 billion click events per day, about 23,000 per second, each roughly 400 bytes as JSON. Raw that is 800 GB a day. Keep four dimensions of moderate cardinality, page (about 50,000 values), country (about 200), device (5) and referrer domain (about 20,000), with one-minute query granularity.

Rollup depends on how many distinct dimension combinations appear per minute. If measurement shows about 300,000 distinct combinations per minute, each minute's 1.4 million events collapse to 300,000 rows, a rollup ratio of about 4.6. That gives about 430 million rows a day. With dictionary encoding and compression, rows like these typically take a few tens of bytes each, so call it 15 to 20 GB of segments per day, roughly 40 to 50 times smaller than the raw JSON. At about 5 million rows per segment, each hourly chunk holds three or four segments, around 90 per day.

Now add a user id dimension. Combinations per minute approach the event count, rollup collapses towards a ratio of 1, and storage grows several-fold. The usual fix is to drop the raw id, keep a sketch metric for distinct users, and send event-level questions to a lake. Measure the ratio on a day of real data before choosing dimensions; it is the single biggest lever on cost.

With 90 days retained on fast Historicals and replication of 2, the hot tier needs roughly 3 TB of segment cache. Historicals perform best when the segments they serve fit in the page cache, so memory, not disk, is usually the real constraint.

Data management: rules, tiers and compaction

Load rules, evaluated in order per datasource, decide where segments live: for example, keep the last 30 days on a hot tier with two replicas, the last year on a cheaper tier with one, and drop older data from Historicals while leaving it in deep storage. A load rule with zero replicas keeps segments queryable from deep storage by the multi-stage engine, slowly, without spending Historical memory on them; a drop rule takes them out of the queryable set. Separate kill tasks delete unused segments from deep storage permanently, and should be configured deliberately because they are irreversible.

Streaming ingestion produces many small segments, especially with late data spread across old time chunks. Compaction rewrites a time chunk into fewer, larger, better-sorted segments with a new version. Run auto-compaction on every streaming datasource, set it to skip the most recent chunks that are still receiving data, and watch the number of segments per interval. Query cost scales with segment count, so compaction lag shows up directly as dashboard latency.

Failure modes

SymptomUsual causeFix
Ingestion tasks never finishHandoff stuck: Historicals full, load rules exclude the interval, or Coordinator not runningCheck Coordinator logs and tier capacity; fix rules
Fresh data missing after a task restartOne replica and the task failed before publishingRun two replicas; tasks re-read from the committed offset
Late events droppedRejection periods configured on the supervisorWiden them or backfill late data in batch
Slow queries on recent dataThousands of small segmentsAuto-compaction with a skip offset; tune task duration
Broker out of memoryLarge groupBy or join merged on one BrokerLimit result sizes, use the multi-stage engine for heavy queries
Storage much larger than plannedHigh-cardinality dimension killed rollupDrop it, use sketches, re-ingest
Ingestion halts cluster-wideMetadata store or deep storage unavailableHistoricals keep serving loaded data; restore the dependency first

When Druid fits and when it does not

Druid fits event data with a strong time component, many concurrent users, and queries that filter and aggregate rather than return raw rows. It fits less well when you need updates to individual records, frequent large joins, or ad hoc exploration over every raw column; a warehouse or lakehouse engine, like the SQL-on-Hadoop engines described in Impala, handles those better. Many organisations run both: the lake holds raw events, Druid holds rolled-up recent data for interactive products. The OLTP versus OLAP distinction is the first filter.

What to do next

  1. Pick one high-traffic dashboard and write down its queries, time ranges and freshness need before designing a datasource.
  2. Load one day of real events with and without each candidate dimension and measure the rollup ratio.
  3. Choose segment granularity from query time ranges and volume, aiming for segments of a few million rows.
  4. Run streaming ingestion with two replicas and enable auto-compaction with a skip offset from day one.
  5. Write load rules for hot, warm and dropped tiers, and decide explicitly whether kill tasks run.
  6. Monitor handoff time, segments per interval, query latency per datasource and Historical cache usage; alert on handoff lag.
Key takeaway: Druid is fast because it does a narrow thing very well: time-partitioned, immutable, columnar segments with bitmap indexes and optional rollup, scanned in parallel and merged on a Broker. Its operating model follows from that: data servers are disposable caches of deep storage, streaming tasks serve fresh data until handoff, and compaction and load rules control cost and latency. Design dimensions for rollup, run replicas and compaction from the start, and protect the metadata store and deep storage like the databases they are.