Druid 把数据存成按时间分片、不可变、列式编码的 Segment。Segment 一旦发布就不能改,但数据还在源源不断地进来,写入后还要能立刻查到。这套”边写边查”的能力是怎么实现的?答案是:所有摄入都以 task 为粒度,由 Overlord 调度。批摄入把存量数据转成 Segment,流摄入让 Kafka 里的数据在到达后几秒内可查。两种模式底层共用同一套 task 执行机制,差异在于任务的生命周期和数据的来源。这篇先讲摄入方式的全景,再聚焦流式摄入的 Supervisor 模型和实时 Segment 的 handoff 机制,最后覆盖 schema 变更与 auto-compaction。
一、摄入方式全景
1.1 一切皆 task
Druid 的官方文档开宗明义:把数据加载进 Druid 的过程叫 ingestion 或 indexing。无论哪种方式,底层都对应一个或多个 indexing task。task 在 Middle Manager 或 Indexer 上执行,读取源数据,构建 Segment,发布到 Deep Storage,再由 Historical 加载服务查询。
这意味着摄入不是一个持续的管道进程,而是一组有生命周期的任务。批摄入的 task 跑完即结束,流摄入的 task 由 Supervisor 持续管理、定期轮换。理解这一点,才能理解 Druid 摄入的调度逻辑:Overlord 是任务调度器,Middle Manager 是任务执行器,task 是最小调度单位。
1.2 批量摄入与流式摄入
Druid 的摄入方式分两大类,下表对照了官方文档列出的几种方法。
| 类别 | 方法 | task 类型 | 适用场景 |
|---|---|---|---|
| 流式 | Kafka 索引服务 | kafka(由 Supervisor 管理) | Kafka 实时数据流 |
| 流式 | Kinesis 索引服务 | kinesis(由 Supervisor 管理) | Amazon Kinesis 实时数据流 |
| 批量 | 原生批处理 | index_parallel | 本地或外部文件批量加载 |
| 批量 | 基于 SQL 的批处理 | MSQ 的 query_controller | 用 INSERT/REPLACE 语句批量加载 |
| 批量 | Hadoop 批处理 | index_hadoop | 已在 37.0.0 移除,改用前两者 |
Hadoop 批处理(index_hadoop)曾是 Druid 经典的批量摄入方式,利用 MapReduce 作业读取 HDFS 上的数据。Apache Druid 37.0.0 起,官方文档明确说该功能已移除,建议用原生批处理或基于 SQL 的批处理替代。Druid 仍支持把 HDFS 作为 Deep Storage,只是不再用 Hadoop MapReduce 做摄入。如果维护的是较老版本的 Druid,index_hadoop 仍然可用,但新项目不要再选它。
流式摄入由一个持续运行的 Supervisor 控制,支持 exactly-once 语义。批量摄入的 task 关联一个 controller task,跑完即结束。两种方式产出的 Segment 没有本质区别,区别在于数据来源和任务的组织方式。
二、批量摄入
2.1 原生批处理:index_parallel
index_parallel 是原生批处理的并行 task 类型。官方文档说它是 supervisor task,负责协调整个摄入过程:先把输入数据拆分,创建 worker task 处理各个数据分片,等所有 worker 成功后统一发布 Segment。
{ "type": "index_parallel", "spec": { "dataSchema": { "dataSource": "wikipedia", "timestampSpec": { "column": "timestamp", "format": "auto" }, "dimensionsSpec": { "dimensions": ["page", "language", "user"] }, "metricsSpec": [ { "type": "count", "name": "count" }, { "type": "doubleSum", "name": "added", "fieldName": "bytes_added" } ], "granularitySpec": { "segmentGranularity": "day", "queryGranularity": "none", "intervals": ["2013-08-31/2013-09-01"] } }, "ioConfig": { "type": "index_parallel", "inputSource": { "type": "local", "baseDir": "examples/indexing/", "filter": "wikipedia_data.json" }, "inputFormat": { "type": "json" } }, "tuningConfig": { "type": "index_parallel", "maxNumConcurrentSubTasks": 2 } }}这个 spec 展示了摄入规范的三个核心组件:dataSchema 定义数据源名、主时间戳、维度和指标;ioConfig 定义数据来源和解析方式;tuningConfig 控制执行参数。所有摄入方式都遵循这个三段式结构,只是各部分的可选值不同。
并行执行靠 tuningConfig 里的 maxNumConcurrentSubTasks 控制。设为大于 1 时,supervisor task 把输入数据拆给多个 worker 并行处理。设为 1 或不设,任务串行执行。官方文档说 worker 失败时 supervisor 会重试,重试次数达到上限才判定整个任务失败。
index_parallel 还有一个简化版:index 类型的单任务摄入,不开并行,适合开发测试。生产环境用 index_parallel。
2.2 perfect rollup 与 best-effort
批量摄入支持两种 rollup 模式。perfect rollup 在摄入时把所有相同维度和时间的行完全合并,保证数据量最小,但需要先扫描一遍数据确定分区,耗时更长。best-effort rollup 不做全局优化,边读边聚合,速度快但合并率可能不如 perfect。
原生批处理默认是 best-effort。要在 index_parallel 里启用 perfect rollup,需要在 tuningConfig 里设 forceGuaranteedRollup = true。基于 SQL 的批处理则始终是 perfect rollup。
2.3 基于 SQL 的批处理
基于 SQL 的摄入是较新的特性,通过多阶段查询(Multi-Stage Query,MSQ)任务引擎执行。用 INSERT 和 REPLACE 语句把数据写进 Druid,底层仍然是 batch task,产出 Segment 的方式和原生批处理一致。
-- 从外部数据源加载到新 datasourceINSERT INTO wikipediaSELECT TIME_PARSE("timestamp") AS __time, "page", "language", "added"FROM TABLE( EXTERN( '{"type": "http", "uris": ["https://druid.apache.org/data/wikipedia.json.gz"]}', '{"type": "json"}', '[{"name": "timestamp", "type": "string"}, {"name": "page", "type": "string"}, {"name": "language", "type": "string"}, {"name": "added", "type": "string"}]' ))PARTITIONED BY DAYEXTERN 函数是 SQL 摄入的关键,它把原生批处理的 inputSource 和 inputFormat 包装成 SQL 可查的表。一条 INSERT/REPLACE 语句至少占用两个 task slot:一个 controller task,至少一个 worker task。
INSERT 和 REPLACE 的区别在于覆盖语义。INSERT 可以创建新 datasource 或追加数据,获取共享锁,多条 INSERT 可并发。REPLACE 用 OVERWRITE 子句指定覆盖范围,可以覆盖整张表或某个时间区间。官方文档提醒:不要用 INSERT 做微批写入,微批场景应该用流式摄入。原因是每条 INSERT 都会生成并发布 Segment,频繁微批会产生大量小 Segment。
三、流式摄入:Kafka 索引服务
3.1 Supervisor 与 indexing task
流式摄入是 Druid 实时性的核心。以 Kafka 为例,官方文档的描述很直接:启用 Kafka 索引服务后,在 Overlord 上配置 Supervisor,由 Supervisor 管理 Kafka indexing task 的创建和生命周期。Kafka indexing task 用 Kafka 的 partition 和 offset 机制保证 exactly-once 摄入。Supervisor 监视这些 task 的状态,协调 handoff,处理失败,维护扩展性和副本需求。
Supervisor 不是直接消费数据的进程,而是一个控制循环。它周期性执行管理逻辑(默认每 30 秒一轮,由 period 控制),根据当前 task 的运行状态决定是否创建新 task、终止失败 task、或推进 handoff。真正读 Kafka 的是它管理的一组 indexing task。
3.2 taskCount、replicas 与分区分配
Supervisor 的 ioConfig 里有几个关键参数控制 task 的并发度。
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
taskCount | 整数 | 1 | 单个副本集里读取 task 的最大数量 |
replicas | 整数 | 1 | 副本集数量,1 表示无副本 |
taskDuration | ISO 8601 | PT1H | task 停止读取并开始发布 Segment 前的运行时长 |
period | ISO 8601 | PT30S | Supervisor 执行管理逻辑的间隔 |
completionTimeout | ISO 8601 | PT30M | 发布超时阈值,超时则终止 task |
taskCount 和 replicas 的乘积是读取 task 的最大数量。官方文档明确说了一个重要约束:当 taskCount 大于 Kafka 分区数时,实际运行的读取 task 数会少于 taskCount。原因是每个 task 至少需要一个分区才有意义,分区不够时多出的 task 不会被创建。
这意味着 task 数量的上限由 Kafka 分区数决定。如果吞吐不够想加 task,得先确认 Kafka topic 的分区数是否够用。replicas 控制副本,Druid 会把同一副本集的 task 分配到不同的 worker 上,单个 worker 挂了不影响摄入。副本之间消费相同的分区,互为热备。
3.3 task 生命周期
一个 indexing task 的生命周期由 taskDuration 控制。默认 PT1H 意味着 task 运行 1 小时后停止读取新数据,开始发布已积累的 Segment。Supervisor 在 task 到期前创建新 task 接替,新 task 从旧 task 停下的 offset 继续消费,实现无缝衔接。
发布阶段不是瞬间完成的。task 停止读取后,要把内存中积累的数据落盘、构建 Segment、写入 Deep Storage、写元数据。如果这个过程超过 completionTimeout(默认 30 分钟),Supervisor 判定 task 失败并终止。这个超时设得太短,Segment 可能还没发完就被杀掉。
四、实时 Segment 生成与 handoff
4.1 内存中的增量构建
indexing task 在 Peon 或 Indexer 的 JVM 里运行,读到的数据先在内存中增量构建 Segment。这个内存中构建 Segment 的组件在 Druid 源码里叫 Appenderator,负责把行数据按列式格式组织、维护字典和位图索引、做 rollup 聚合。官方文档在描述 task 行为时用”构建 Segment”概括这个过程,Appenderator 是其内部实现名称。
内存不是无限的。数据积累到 maxRowsInMemory(默认 15 万行,指聚合后的行数)或 maxBytesInMemory(默认为 JVM 堆的六分之一)时,触发一次 intermediate persist(中间持久化),把内存数据写到本地磁盘的临时文件里,腾出内存继续接收新数据。这个过程类似 LSM 树的 flush。
intermediatePersistPeriod(默认 PT10M)控制中间持久化的时间间隔。即使行数没达到 maxRowsInMemory,超过这个时间也会触发一次 persist。maxPendingPersists 限制排队等待的 persist 数量,超过则阻塞摄入。官方文档给出的堆内存占用估算公式是 maxRowsInMemory * (2 + maxPendingPersists),配置时要注意。
4.2 handoff 的触发条件
handoff 指 Segment 从 indexing task 交接给 Historical 服务的过程。官方文档在 tuningConfig 说明里明确了三个触发条件,满足任一即触发:
| 条件 | 默认值 | 说明 |
|---|---|---|
maxRowsPerSegment | 5000000 | 单个 Segment 的行数上限(聚合后行数) |
maxTotalRows | 20000000 | 所有 Segment 累计行数上限 |
intermediateHandoffPeriod | P2147483647D | 定期 handoff 的时间间隔 |
默认的 intermediateHandoffPeriod 是一个极大的天数(P2147483647D,约 588 万年),实际上等于”永不按时间触发”。这意味着默认情况下 handoff 完全由行数驱动。如果需要按时间周期性 handoff,得显式设置这个参数。
maxRowsPerSegment 和 maxTotalRows 指的都是 post-aggregation rows,即 rollup 聚合后的行数,不是原始输入事件的条数。如果 rollup 把行数压了一个数量级,500 万行的 Segment 可能对应几千万条原始事件。
4.3 handoff 的过程
Kafka indexing task 支持增量 handoff(incremental hand-off)。官方文档说,当 task 达到上述任一条件时,它不会等到 taskDuration 结束,而是立即把当前积累的 Segment 发布出去,然后创建新的 Segment 继续接收后续数据。这让 Segment 在生成后很快可用,不必等整个 task 周期结束。
发布一个 Segment 的完整步骤:合并本地临时持久化的文件,构建完整的列式 Segment 文件,写入 Deep Storage,把 Segment 的元数据(标识符、时间区间、版本、Deep Storage 位置)写入 Metadata Storage。Coordinator 轮询 Metadata Storage 发现新 Segment,指示 Historical 从 Deep Storage 下载并加载。加载完成后,该 Segment 由 Historical 服务查询。
在 Historical 加载完成前,查询这个时间区间的数据会同时发给 Historical 和 Middle Manager。Middle Manager 上的实时 Segment 直接服务查询,Historical 上的已发布 Segment 服务历史数据。Broker 负责合并两边的结果。这就是 Druid “写入后立即可查”的具体机制:数据在 indexing task 的内存里就能被查到,不需要等 handoff 完成。
handoff 期间会短暂出现同一个 Segment 的两个版本:indexing task 上的实时版本和 Historical 上的已发布版本。Druid 用 Segment 版本号管理这个过渡,新版本覆盖旧版本后,indexing task 上的旧版本被清除。
五、延迟数据处理
流式摄入的 Kafka 数据不一定按时间戳顺序到达。一条事件可能因为网络延迟、生产端重试等原因,在生成很久之后才被消费。Supervisor 的 ioConfig 提供两个参数处理这种 late data(延迟数据)。
lateMessageRejectionStartDateTime 指定一个绝对时间点,早于这个时间戳的消息会被丢弃。lateMessageRejectionPeriod 指定一个相对时长,早于 task 创建时间减去这个时长的消息会被丢弃。两者只能设一个。比如设 lateMessageRejectionPeriod 为 PT1H,Supervisor 在 12<00>00> 创建的 task 会丢弃时间戳早于 11<00>00> 的消息。
earlyMessageRejectionPeriod 处理另一个方向的异常:时间戳太靠未来的消息。如果一个消息的时间戳超过 task 到达 taskDuration 后再加上这个时长,也会被丢弃。这是防止脏数据或时钟漂移把 Segment 的时间范围撑得离谱。
延迟数据之所以需要处理,官方文档给出的场景是:如果你的数据流有延迟消息,且同时有流式和批量两条摄入管道操作相同的 Segment,延迟消息可能触发并发问题。丢弃过旧的消息可以避免实时管道和夜间批量管道在同一个时间区间上冲突。
六、schema 变更
6.1 新数据:直接改 spec
Druid 的 Segment 在创建时存储一份自己的 schema 副本。官方文档说,这意味着对新写入的数据,schema 变更不需要动已有数据。流式摄入更新 Supervisor spec,下次批量摄入用新 schema 即可。Druid 在查询时自动协调不同 Segment 的 schema 差异。
比如给 DataSource 加一个新维度列,只需要在 Supervisor spec 的 dimensionsSpec 里加上这个字段。Supervisor 轮换 task 后,新 task 产出的 Segment 会包含新列。旧 Segment 没有这个列,查询时 Druid 把缺失值当 null 处理。这种前向兼容让加列操作很轻量。
6.2 已有数据:reindex
要修改已有数据的 schema,比如改列的类型或删列,就不能只改 spec 了。官方文档说这需要 reindexing(重新摄入),即读取所有受影响的 Segment,用新 schema 重新生成并发布。这与 Segment 不可变的特性一致:数据一旦发布就不能原地改,只能用新版本覆盖。
reindex 是个耗时操作,因为它要重写所有涉及的 Segment。auto-compaction(下一节)本质上也是一种 reindex,只是它的目的是优化 Segment 大小而不是改 schema。两者可以结合:在 compaction 任务里同时修改 schema,一步完成优化和变更。
6.3 Supervisor 的 schema 变更策略
流式摄入改 schema 时,Supervisor 需要切换到新 spec。最直接的方式是 suspend 当前 Supervisor,提交新 spec 的 Supervisor,让它从之前的 offset 继续。这个过程会有短暂的摄入中断。
Druid 也支持更平滑的方式:提交新 spec 后,Supervisor 在下一次 task 轮换时用新 schema 创建 task。旧 task 运行到 taskDuration 结束后自然退出,新 task 接替。这样新旧 Segment 在过渡期共存,Druid 查询时自动协调。具体行为取决于 Supervisor 版本和配置,不确定时参考官方文档对应版本的说明。
七、auto-compaction
7.1 为什么需要 compaction
流式摄入有一个固有的问题:Segment 按 taskDuration 周期生成,每个周期产出一批 Segment。如果 taskDuration 设为 1 小时、Kafka 有 4 个分区,每小时就产出 4 个 Segment(每个 task 一个)。一天下来 96 个 Segment,每个可能只有几十 MB。这些小 Segment 增加查询时的调度开销,也占用更多 Metadata Storage 的记录。
auto-compaction 解决的就是这个问题。官方文档说,compaction task 读取一个时间区间的现有 Segment,把数据合并成新的”压缩后”的 Segment。压缩后的 Segment 通常更大但数量更少,查询时需要的 per-segment 处理开销和内存开销都更低。
7.2 Coordinator 触发
auto-compaction 由 Coordinator 触发。官方文档描述它的行为:Coordinator 用 segment search policy 周期性识别需要 compaction 的 Segment,从最新到最旧扫描。发现未被压缩过、或用过时 spec 压缩过的 Segment 后,提交 compaction task 处理对应的时间区间。
auto-compaction 有两种运行方式。较新的是 compaction supervisor(推荐),在 Overlord 上用 Supervisor 框架管理,支持 MSQ 任务引擎,响应更快,能通过 supervisor API 查看状态。传统方式是作为 Coordinator 的 duty 运行。两者效果类似,管理方式不同。
compaction 的配置是动态的,不需要重启 Druid。核心字段包括 dataSource、taskPriority、skipOffsetFromLatest(避开最新数据,减少与摄入任务的冲突)、以及可选的 granularitySpec、dimensionsSpec、metricsSpec(用于在 compaction 时调整 schema)。
auto-compaction 跳过 segmentGranularity 为 ALL 的 datasource。ALL 意味着整个 datasource 是一个时间块,没有细分,compaction 对它无意义。
7.3 compaction 与摄入的冲突
compaction task 和摄入 task 可能操作同一个时间区间。官方文档说,默认情况下摄入 task 优先:如果一个摄入 task 要往某个正在被 compaction 的时间区间写数据,compaction task 会失败退出。可以通过调高 compaction task 的优先级来改变这个行为,或者用 skipOffsetFromLatest 让 compaction 避开最新时间区间。
这也是为什么 auto-compaction 通常只处理较旧的数据:最新数据还在被实时 task 写入,compaction 碰它容易冲突。等数据”成熟”一段时间后,实时 task 已经 handoff 给 Historical,compaction 再介入就安全了。
八、踩坑
8.1 taskCount 超过分区数
前面讲过,taskCount 大于 Kafka 分区数时多余的 task 不会创建。这个限制容易被忽略:有人觉得加 taskCount 就能提升吞吐,但 Kafka topic 只有 2 个分区时,taskCount 设 4 也只有 2 个 task 在跑。要提升并行度,得先给 Kafka topic 加分区。加分区只对新数据生效,已有分区的数据不会重新分配。
8.2 小 Segment 爆炸
taskDuration 设得太短、segmentGranularity 设得太细,或者 Kafka 分区数多但每个分区数据量小,都会产生大量小 Segment。官方文档给的例子:taskDuration 4 小时、segmentGranularity 1 小时、Supervisor 9<10>10> 启动,13<10>10> 换新 task 时,13<00>00> 到 14<00>00> 的数据会被旧 task 和新 task 同时写入,产生跨 task 的小 Segment。
这些小 Segment 需要靠 auto-compaction 在后台合并成 300 到 700 MB 的理想大小。如果 compaction 跟不上,Segment 数量持续膨胀,查询性能下降,Metadata Storage 的负载也上升。监控 Segment 数量是 Druid 运维的基本功。
8.3 handoff 卡住
handoff 依赖 Historical 加载 Segment。如果 Historical 节点磁盘满了、或者 Coordinator 因为负载规则没有把 Segment 分配出去,handoff 会一直卡着。indexing task 的 Segment 发布了但 Historical 迟迟不加载,查询会继续走 Middle Manager,占用摄入任务的资源。
排查方向是看 Coordinator 日志里 Segment 的分配情况,确认 Historical 有足够的可用磁盘空间和 task slot。segmentAvailabilityConfirmed 这个 task report 字段能告诉你 task 产出的 Segment 是否已被集群确认可用。
8.4 延迟数据被静默丢弃
lateMessageRejectionPeriod 配置后,早于阈值的消息会被丢弃,而且默认没有明显的告警。如果上游有大量延迟消息(比如重放历史数据),配了 PT1H 的拒绝窗口,所有超过 1 小时的消息都悄无声息地消失。排查时看 task report 里的 thrownAway 计数,它记录了被时间窗口拒绝的消息数。
如果需要重放历史数据,临时把 lateMessageRejectionPeriod 调大或去掉,重放完再恢复。或者用 useEarliestOffset 让新 Supervisor 从最早 offset 开始消费,但这会消费全量数据,慎用。
参考资料
- Apache Druid - Ingestion overview - 摄入方式总览,批量与流式的分类对照
- Apache Druid - Native batch ingestion - index_parallel 任务类型、并行执行、append 与 overwrite 语义
- Apache Druid - Hadoop-based ingestion - 37.0.0 移除说明,建议改用原生批处理或 SQL 摄入
- Apache Druid - SQL-based ingestion concepts - MSQ 任务引擎、INSERT/REPLACE 语义、EXTERN 函数
- Apache Druid - Apache Kafka ingestion - Kafka 索引服务配置、分区分配、增量 handoff
- Apache Druid - Supervisor - Supervisor spec、taskCount/replicas/taskDuration 参数、tuningConfig
- Apache Druid - Task reference - task API、task report、Segment 可用性字段
- Apache Druid - Ingestion spec reference - dataSchema、ioConfig、tuningConfig 三段式结构
- Apache Druid - Compaction - compaction 机制、Segment 与 query granularity 处理
- Apache Druid - Automatic compaction - auto-compaction 配置、compaction supervisor、skipOffsetFromLatest
- Apache Druid - Schema changes - 新数据与已有数据的 schema 变更策略、reindexing
支持与分享
如果这篇文章对你有帮助,欢迎支持作者或分享给更多人
部分信息可能已经过时






