资讯动态

Apache Druid Indexing Service 架构深度解析:Overlord、MiddleManager 与 Peon 的任务编排机制

发布时间:2026/9/23 1:54:23 来源:尧图企业网站定制
Apache Druid Indexing Service 架构深度解析Overlord、MiddleManager 与 Peon 的任务编排机制【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址: https://gitcode.com/gh_mirrors/druid6/druidApache Druid 的 Indexing Service索引服务是支撑全部分布式数据摄取与段管理能力的高可用服务层。本文以仓库中 docs/design/indexing-service.md 为骨架系统讲解 Druid 如何通过 Overlord、MiddleManager、Peon 三级组件完成任务的接收、调度、隔离执行与结果回报并深入对应的源码实现indexing-service 模块与任务 APItasks.md验证其底层原理。读完本文你将掌握 Indexing Service 的部署拓扑、Overlord 两种运行模式、任务状态机、锁与优先级机制以及日志与报告的实际配置方法。1. Indexing Service 是什么职责与总体架构Apache Druid 的 Indexing Service 是一个高可用、分布式的服务专门负责运行与索引Indexing相关的各类任务。它承担着 Druid 数据写入链路中的核心编排职责所有创建 Segments段以及 Kill 段 的工作都以“任务Task”的形式提交给该服务执行。换句话说Druid 的批量摄取Batch Ingestion、流式摄取Streaming Ingestion、压缩Compaction与数据清理Kill等能力全部由 Indexing Service 承载。从仓库中的 docs/design/indexing-service.md 可以看到Indexing Service 由三个核心组件组成它们构成了一个清晰的“调度—管理—执行”三级模型组件职责运行位置Peon单个任务的执行引擎一个 Peon 只能运行一个任务与 MiddleManager 同机MiddleManager管理一组 Peon负责任务的转发与资源隔离与 Peon 同进程部署Overlord负责任务分发Task Distribution将任务分配给 MiddleManager可与 MiddleManager 同进程或独立部署需要特别强调的是部署上的约束关系Overlord 与 MiddleManager既可以运行在同一进程中也可以拆分到多个进程甚至不同服务器上MiddleManager 与 Peon 则始终运行在同一进程中因为 Peon 是由 MiddleManager 直接派生spawn出来的子 JVM。任务的提交与管理统一通过 Overlord 服务暴露的 HTTP API 完成详细接口清单见 Tasks API 参考。下图展示了 Indexing Service 的总体组件关系来源docs/assets/indexing_service.png1.1 从源码看 Indexing Service 的进程边界上述“组件—进程”关系在仓库的 services 模块中有直接印证。三个组件分别对应三个 CLI 入口CliOverlord.javaCliMiddleManager.javaCliPeon.java它们都挂在统一的启动类org.apache.druid.cli.Main下。MiddleManager 的启动方式为见 docs/design/middlemanager.mdorg.apache.druid.cli.Main server middleManager而 Peon 的启动方式比较特殊它需要显式传入任务文件与状态文件见 docs/design/peons.mdorg.apache.druid.cli.Main internal peon task_file status_file其中task_file是任务 JSON 对象文件status_file指定任务状态输出位置。从该命令可以看出Peon 是被“按任务”拉起的一次性执行进程——这正印证了文档中“Peons should rarely run on their ownPeon 几乎不应单独运行”的描述Peon 仅用于开发调试场景。2. Overlord任务调度的中枢Overlord 服务是整个 Indexing Service 的控制面它负责接收任务accepting tasks对外提供任务提交 API协调任务分发coordinating task distribution决定把任务交给哪个 MiddleManager创建任务锁creating locks around tasks保证同一数据源同一时间块内数据写入的正确性返回状态给调用方returning statuses to callers供客户端查询任务进度。其源码骨架位于 indexing-service/src/main/java/org/apache/druid/indexing/overlord/核心协作类包括TaskQueue.java任务生产者与 TaskRunner 之间的接口负责接收任务、按Task.isReady判定就绪状态、以近似 FIFO 顺序交付给 TaskRunner并将任务与状态变更持久化到 TaskStorageTaskMaster.java封装 Overlord 主节点Leader的各类状态类提供对 TaskRunner、TaskQueue 的访问TaskRunner.java任务运行器抽象其下派生 Local本地与 Remote远程两种实现。2.1 本地模式Local与远程模式Remote根据 docs/design/overlord.mdOverlord 支持两种运行模式默认是本地模式Local Mode本地模式Overlord 同时负责创建执行任务的 Peon。此时必须一并提供 MiddleManager 与 Peon 的全部配置。该模式适合简单工作流相当于把 Overlord 与 MiddleManager 合并在同一个进程中运行。远程模式Overlord 与 MiddleManager 作为独立服务运行可以分别部署在不同服务器上。文档明确建议如果你打算把 Indexing Service 作为所有 Druid 索引操作的唯一入口应使用远程模式。两种模式的差异本质上对应 TaskRunner 的两类实现本地任务运行器ForkingTaskRunner / ThreadingTaskRunner 等直接在当前进程内 fork 子 JVM远程任务运行器 RemoteTaskRunner.java 则通过 ZooKeeper 发现集群中的 MiddleManager 并将任务转发过去。2.2 Overlord 的配置与调优入口Overlord 服务的完整配置参数位于 Overlord 配置基础调优建议参见 Basic cluster tuning。2.3 Worker 黑名单机制Blacklisted Workers当一个 MiddleManager 的任务失败次数超过阈值时Overlord 会将其加入黑名单blacklist不再向它分配新任务。黑名单机制有两个关键约束最多 20% 的 MiddleManager 可以被列入黑名单由druid.indexer.runner.maxPercentageBlacklistWorkers控制被列入黑名单的 MiddleManager 会周期性自动恢复whitelist。控制该行为的四个配置项如下见 docs/design/overlord.mddruid.indexer.runner.maxRetriesBeforeBlacklist druid.indexer.runner.workerBlackListBackoffTime druid.indexer.runner.workerBlackListCleanupPeriod druid.indexer.runner.maxPercentageBlacklistWorkers这些配置在源码 RemoteTaskRunnerConfig.java 中有对应实现与默认值例如maxRetriesBeforeBlacklist默认5次workerBlackListBackoffTime默认PT15M15 分钟退避时间。2.4 自动伸缩AutoscalingOverlord 还内置了自动伸缩框架Autoscaling。当前实现与 Druid 官方自有的部署基础设施耦合较紧历史上 MiddleManager 以 AWS EC2 节点形式运行并注册到 Galaxy 环境但框架本身预留了扩展点官方欢迎社区提供新的实现。自动伸缩的触发条件有两个方向扩容当某个任务处于 Pending等待状态时间过长时可以新增 MiddleManager缩容当 MiddleManager 在一段时间内未运行任何任务时可以被终止。2.5 Overlord 的 HTTP 端点Overlord 暴露的 API 端点清单见 Service status API reference。最常用的任务管理接口提交、取消、查询状态、查看日志与报告则在 Tasks API reference 中有完整说明。3. MiddleManager任务执行的工作节点MiddleManager 是执行已提交任务的 Worker 服务。它本身并不直接执行任务而是把任务转发给运行在独立 JVM中的 Peon。3.1 为什么任务要跑在独立 JVM 中根据 docs/design/middlemanager.mdDruid 使用独立 JVM 运行任务的核心目的是隔离资源与日志每个 Peon 一次只能运行一个任务一个 MiddleManager 可以同时管理多个 Peon受druid.worker.capacity槽位数限制单个任务崩溃、OOM 或日志爆炸不会拖垮整个 MiddleManager 及其他任务。3.2 MiddleManager 的配置与运行MiddleManager 与 Peon 的配置集中在 MiddleManager and Peons 配置基础调优指导见 Basic cluster tuning。其启动命令见 docs/design/middlemanager.mdorg.apache.druid.cli.Main server middleManagerHTTP 端点清单见 Service status API reference。4. Peon单任务执行引擎Peon 是由 MiddleManager 派生的任务执行引擎每个 Peon 运行一个独立 JVM一个 Peon 只负责执行单个任务Peon 永远运行在派生它的 MiddleManager 所在主机上见 docs/design/peons.md。4.1 Peon 的配置与运行Peon 的配置分为两部分Peon Query Configuration 与 Additional Peon Configuration。任务级调优参见 Basic cluster tuning。正常情况下 Peon 不应被单独运行仅用于开发目的命令形式为org.apache.druid.cli.Main internal peon task_file status_file4.2 任务槽位与任务存储大小的分配逻辑任务在执行期间可能需要本地磁盘例如实时摄取任务接收广播段、Multi-stage Query 的中间数据集。根据 docs/ingestion/tasks.md 的“Configuring task storage sizes”一节任务存储大小由三个配置项共同决定druid.worker.capacity即“任务槽位数”druid.worker.baseTaskDirs用于任务存储的目录列表druid.worker.baseTaskDirSize每个存储位置上可用的存储量。分配逻辑有两个要点一个任务只使用一个目录虽然配置了目录列表但任意给定任务只会从列表中选取一个目录作为临时scratch空间不会同时占用多个按槽位均分容量每个任务实际分到的磁盘空间是“能让所有任务槽位获得等量磁盘”的最大值。文档给出的示例是5 个槽位、2 个目录A 和 B、每个目录 300 GB 时目录 A 分到 3 个槽位、目录 B 分到 2 个槽位每个槽位允许使用 100 GB。5. TasksIndexing Service 的“工作单元”任务Task是 Indexing Service 处理的对象。根据 docs/ingestion/tasks.md批量摄取场景下一般直接通过 Tasks API 向 Druid 提交任务流式摄取场景下任务通常由 Supervisor如 Kafka/Kinesis Indexing Service代为自动提交。5.1 任务 API 的两个入口任务 API 主要有两个使用位置Overlord 进程提供提交、取消、查询状态、查看日志与报告等 HTTP API完整清单见 Tasks API referenceDruid SQL 的sys.tasks表见 SQL metadata tables提供当前运行中任务的只读信息虽然字段比 Overlord API 少但对日常巡检非常实用。5.2 任务报告Task Reports任务报告包含摄取行数与解析异常信息对已完成和运行中的任务都可获取。支持报告功能的包括 native batch 任务、Hadoop 批量任务以及 Kafka、Kinesis 摄取任务。完成报告Completion Report任务完成后可通过如下端点获取报告http://OVERLORD-HOST:OVERLORD-PORT/druid/indexer/v1/task/{taskId}/reports一个典型输出如下原文档示例字段含义见下文{ ingestionStatsAndErrors: { taskId: compact_twitter_2018-09-24T18:24:23.920Z, payload: { ingestionState: COMPLETED, unparseableEvents: {}, rowStats: { determinePartitions: { processed: 0, processedBytes: 0, processedWithError: 0, thrownAway: 0, unparseable: 0 }, buildSegments: { processed: 5390324, processedBytes: 5109573212, processedWithError: 0, thrownAway: 0, unparseable: 0 } }, segmentAvailabilityConfirmed: false, segmentAvailabilityWaitTimeMs: 0, recordsProcessed: { partition-a: 5789 }, errorMsg: null }, type: ingestionStatsAndErrors }, taskContext: { type: taskContext, taskId: compact_twitter_2018-09-24T18:24:23.920Z, payload: { forceTimeChunkLock: true, useLineageBasedSegmentAllocation: true } } }Compaction 任务的特殊性压缩任务可能按输入 interval 的拆分方式生成多组段输出报告因此整体报告会包含从每个拆分split到对应报告的映射例如ingestionStatsAndErrors_0、ingestionStatsAndErrors_1分别对应不同日期分片的结果。段可用性字段Segment Availability Fields字段描述segmentAvailabilityConfirmed该摄取任务生成的所有段在任务完成前是否已被集群确认为可查询segmentAvailabilityWaitTimeMs摄取完成后任务等待新段可查询所花费的毫秒数recordsProcessed摄取任务处理的分区及每个分区处理的记录数Compaction 任务段信息字段字段描述segmentsReadCompaction 任务读取的段数量segmentsPublishedCompaction 任务发布的段数量实时报告Live Report任务运行期间通过同一端点可以获取包含摄取状态、不可解析事件以及1 分钟、5 分钟、15 分钟处理事件数移动平均值的实时报告http://OVERLORD-HOST:OVERLORD-PORT/druid/indexer/v1/task/{taskId}/reports实时报告示例{ ingestionStatsAndErrors: { taskId: compact_twitter_2018-09-24T18:24:23.920Z, payload: { ingestionState: RUNNING, unparseableEvents: {}, rowStats: { movingAverages: { buildSegments: { 5m: { processed: 3.392158326408501, processedBytes: 627.5492903856, unparseable: 0, thrownAway: 0, processedWithError: 0 }, 15m: { processed: 1.736165476881023, processedBytes: 321.1906130223, unparseable: 0, thrownAway: 0, processedWithError: 0 }, 1m: { processed: 4.206417693750045, processedBytes: 778.1872733438, unparseable: 0, thrownAway: 0, processedWithError: 0 } } }, totals: { buildSegments: { processed: 1994, processedBytes: 3425110, processedWithError: 0, thrownAway: 0, unparseable: 0 } } }, errorMsg: null }, type: ingestionStatsAndErrors } }5.3 报告字段语义详解ingestionState任务到达的摄取阶段。可能取值包括NOT_STARTED任务尚未开始读取任何行DETERMINE_PARTITIONS任务正在处理行以确定分区方式BUILD_SEGMENTS任务正在处理行以构建段COMPLETED任务已完成工作。注意只有批量任务有 DETERMINE_PARTITIONS 阶段实时任务如 Kafka Indexing Service 创建的任务没有该阶段。unparseableEvents由不可解析输入导致的异常消息列表DETERMINE_PARTITIONS 与 BUILD_SEGMENTS 阶段各一份。可用于定位问题输入行。Hadoop 批量任务不支持保存不可解析事件。rowStats中的各项行计数processed成功摄取且无解析错误的行数processedBytes任务处理的总未压缩字节数包含processedWithError、unparseable、thrownAway中的行processedWithError被摄取但某列存在解析错误的行数例如数值列收到了非数值字符串thrownAway被跳过的行数包括时间戳超出任务时间区间、以及被 transformSpec 过滤的行不包括由显式配置如 CSV 的skipHeaderRows、hasHeaderRow跳过的行unparseable完全无法解析而被丢弃的行数例如 JSON 解析器收到非 JSON 数据。errorMsg导致任务失败的错误描述任务成功时为 null。5.4 运行中任务的实时行统计与不可解析事件native batch 任务、Hadoop 批量任务以及 Kafka、Kinesis 摄取任务在运行期间支持行统计检索通过运行该任务的 Peon 上的 GET 请求获取http://middlemanager-host:worker-port/druid/worker/v1/chat/{taskId}/rowStats返回结构中的movingAverages是四个行计数器的 1/5/15 分钟移动平均语义与完成报告相同totals是当前累计值。对于 Kafka Indexing Service可通过 Overlord API 汇总该 Supervisor 下所有任务的状态http://OVERLORD-HOST:OVERLORD-PORT/druid/indexer/v1/supervisor/{supervisorId}/stats此外运行中任务的最近不可解析事件列表可通过 Peon API 获取http://middlemanager-host:worker-port/druid/worker/v1/chat/{taskId}/unparseableEvents注意该能力并非所有任务类型都支持目前仅支持非并行 native batch 任务类型index以及 Kafka、Kinesis Indexing Service 创建的任务。6. 任务锁系统与段版本数据正确性的保障Druid 的锁系统与段版本系统紧密耦合共同保证摄取数据的正确性。这一部分在 docs/ingestion/tasks.md 中有完整阐述。6.1 段的“遮蔽Overshadowing”关系运行覆盖任务时新段会**遮蔽overshadow**旧段。遮蔽关系仅在同一数据源datasource的同一时间块time chunk内成立被遮蔽的段不会参与查询处理从而过滤掉过期数据。每个段拥有主版本major version与次版本minor version主版本是yyyy-MM-ddThh:mm:ss格式的时间戳次版本是整数。段s1遮蔽s2的条件是s1的主版本高于s2或s1与s2主版本相同且s1的次版本更高。举例主版本2019-01-01T00:00:00.000Z、次版本0的段遮蔽主版本2018-01-01T00:00:00.000Z、次版本1的段主版本2019-01-01T00:00:00.000Z、次版本1的段遮蔽主版本同为2019-01-01T00:00:00.000Z、次版本0的段。6.2 两类锁时间块锁Time Chunk Lock与段锁Segment Lock如果两个或多个任务为同一数据源的同一时间块生成段生成的段可能相互遮蔽导致错误查询结果。为避免该问题任务在创建任何段之前会尝试获取锁。锁分为两类时间块锁Time Chunk Lock任务锁定整个时间块。持有锁期间其他任务无法为同一数据源的同一时间块创建段。使用时间块锁创建的段其主版本高于已有段次版本恒为0。段锁Segment Lock任务只锁定单个段。因此只要读取的是不同段两个或多个任务可以同时为同一数据源的同一时间块创建段。例如 Kafka 摄取任务与压缩任务总是可以同时写同一时间块——因为 Kafka 摄取任务总是追加新段而压缩任务总是覆盖已有段。使用段锁创建的段主版本相同、次版本更高。注意段锁仍处于实验阶段可能存在导致错误查询结果的未知缺陷。要启用段锁可在任务上下文task context中设置forceTimeChunkLock为false一旦取消设置任务会自动选择合适的锁类型。需要说明的限制当覆盖任务改变段粒度时通常强制使用时间块锁段锁仅受 native 摄取任务与 Kafka/Kinesis 摄取任务支持Hadoop 摄取任务不支持。forceTimeChunkLock只作用于单个任务若要全局关闭可在 Overlord 配置中设置druid.indexer.tasklock.forceTimeChunkLock为false见 Overlord operations 配置。6.3 锁冲突、抢占与锁优先级锁请求在以下情况会冲突两个或多个任务试图为同一数据源的重叠时间块获取锁且冲突可以发生在不同类型的锁之间。冲突行为取决于任务优先级若冲突任务优先级相同先请求者先得锁其他任务等待其释放低优先级任务后于高优先级任务请求锁时低优先级任务等待高优先级任务后于低优先级任务请求锁时高优先级任务会**抢占preempt**低优先级任务的锁低优先级任务的锁被撤销高优先级任务获得新锁。抢占可能发生在任务运行期间的任何时刻发布段publishing segments的关键区critical section除外——发布结束后锁重新变为可抢占。同一 groupId 的任务共享锁。例如同一 Supervisor 下的 Kafka 摄取任务具有相同 groupId彼此共享全部锁。默认锁优先级数值越大优先级越高任务类型默认优先级Realtime index task75Batch index tasks含 native batch、SQL、Hadoop50Merge/Append/Compaction task25Other tasks0可通过任务上下文覆盖优先级context : { priority : 100 }6.4 任务动作Task Actions任务动作是任务生命周期中由 Overlord 执行的操作典型动作包括lockAcquire为任务获取某个时间区间的时间块锁lockRelease释放任务在某个时间区间上获取的锁segmentTransactionalInsert以单个事务发布任务创建的新段并可覆盖/删除已有段segmentAllocate为任务分配 pending segments 以写入行。批量segmentAllocate动作在多个任务并发的集群中Overlord 上的segmentAllocate动作可能耗时很长导致task/action/run/time尖峰进而引发摄取延迟累积。根因通常是多个并发任务为同一数据源和同一时间区间分配段对 segments 与 pending segments 元数据表的大量元数据调用获取分配段所需任务锁时的并发限制。由于争用通常来自同一数据源与时间区间的段分配可通过批量batching动作改善运行时间。在 Overlord 配置中设置druid.indexer.tasklock.batchSegmentAllocation为true即可启用详见 Overlord operations 配置。6.5 任务上下文参数Context Parameters任务上下文用于各种单任务配置在摄取规范ingestion spec的context字段中指定。配置自动压缩时应把上下文配置放在taskContext而非context中——这些设置会被传入下发给 MiddleManager 的压缩任务的context字段。适用于所有任务类型的参数如下属性描述默认值forceTimeChunkLock设为 false 仍属实验特性。强制使用时间块锁。为true时覆盖 Overlord 运行属性druid.indexer.tasklock.forceTimeChunkLock。两者均非true时各任务自动选择锁类型。truepriority任务优先级取决于任务类型storeCompactionState是否在元数据存储中保存所创建段的压缩状态。为true时任务创建的段会填充段元数据中的lastCompactionState。压缩任务自动设置该参数。压缩任务为true其他任务为falsestoreEmptyColumns摄取时是否存储空列。为true时存储 dimensionsSpec 中指定的每一列为false时Druid SQL 查询引用空列会失败。若保持禁用应摄入占位数据或避免查询空列。该参数覆盖系统属性 druid.indexer.task.storeEmptyColumns。truetaskLockTimeout任务锁超时毫秒。任务获取锁时通过 HTTP 发送请求并等待锁获取结果因此若taskLockTimeout大于 Overlord 的druid.server.http.maxIdleTime可能出现 HTTP 超时错误。300000useLineageBasedSegmentAllocation为动态分区的 native Parallel 任务启用新的基于血缘的段分配协议。在从 0.190.21 版本滚动升级到 0.22 或更高版本期间应关闭升级完成后必须设为true以保证数据正确性。0.21 及更早为false0.22 及更晚为true7. 任务日志本地目录与长期存储任务运行时会创建日志。根据 docs/ingestion/tasks.md 的“Task logs”一节日志流程如下任务提交给 Overlord 后先处于WAITING状态等待获取锁Worker 槽位分配后进入PENDING状态直到任务真正开始执行任务开始在 MiddleManager或 Indexer本地目录的log目录下、以具体taskId为子目录写日志位置由 druid.worker.baseTaskDirs 决定任务完成无论成功失败MiddleManager或 Indexer将日志文件推送到 druid.indexer.logs 指定的位置。Druid Web Console 中的任务日志通过 Overlord 上的 API 获取Overlord 自动探测日志文件所在位置MiddleManager/Indexer 本地或长期存储并返回给前端。排障提示如果在长期存储中看不到日志文件说明要么 (1) MiddleManager/Indexer 推送日志到深存储失败要么 (2) 任务没有完成。可以先检查 MiddleManager/Indexer 本地日志是否有推送失败若无则检查 Overlord 自身进程日志看任务在启动前为何失败。远程模式注意如果以远程模式运行 Indexing Service任务日志必须存储在 S3、Azure Blob Store、Google Cloud Storage 或 HDFS 中。可通过设置druid.indexer.logs.kill系列属性毫秒配置日志保留期见 Task logging 配置Overlord 会自动管理日志目录中的任务日志及任务相关元数据存储表条目。时钟同步注意日志自动删除通常基于后端存储中日志文件的“修改时间戳”Druid 进程与长期存储之间的较大时钟偏移可能导致意外行为。8. Indexing Service 中的任务类型一览Indexing Service 支持的任务类型在 docs/ingestion/tasks.md 的“All task types”一节中列出任务类型说明详细文档index_parallel原生并行批量摄取Native batch ingestion (parallel task)index_hadoopHadoop 批量摄取Hadoop-based ingestionindex_kafkaKafka 流式摄取由 Supervisor 自动提交Kafka-based ingestionindex_kinesisKinesis 流式摄取由 Supervisor 自动提交Kinesis-based ingestioncompact合并指定时间区间的所有段Compactionkill删除指定段的元数据并移除深存储中的数据Deleting data9. 延伸阅读围绕 Indexing Service仓库中还有以下高价值文档与代码可供深入Indexing Service 组件文档总览Overlord 服务文档MiddleManager 服务文档Peon 服务文档Segments 段模型Tasks API 参考Service status API 参考索引服务源码模块服务进程 CLI 入口CliOverlord / CliMiddleManager / CliPeon【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址: https://gitcode.com/gh_mirrors/druid6/druid创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价