深入 Kibana(十二):Task Manager——Kibana 内部的分布式任务调度
上一篇展示了 Reporting 作为 Task Manager 的具体用户。这一篇进入 Task Manager 本身,核心问题是:告警、上报、后台任务共用的调度器;任务如何在多 Kibana 实例间分配;与 cron 的区别。
Task Manager 的定位
Task Manager 是 Kibana 进程内的分布式任务调度器,7.4 版本引入,解决的问题是:Kibana 需要运行定期或一次性后台任务(告警检查、报表生成、ML 模型同步等),但不依赖外部消息队列或 cron 守护进程。
1 | |
所有需要调度的 Kibana 功能——Alerting、Reporting、ML、Fleet——都是 Task Manager 的消费者,注册自己的 task type,由 Task Manager 统一调度。
任务存储:.kibana_task_manager 索引
每个任务对应索引中的一条文档,关键字段如下:
1 | |
ownerId 是 claim 机制的核心:null 表示任务空闲可被抢占,非 null 表示被某个 Kibana 实例持有。
claim 机制与乐观锁
1 | |
乐观锁通过 Elasticsearch 的 if_seq_no + if_primary_term 实现,无需分布式锁服务。多个实例同时 claim 同一任务时,只有一个能成功更新,其余得到 409 conflict 后跳过。
循环任务与一次性任务
| 属性 | 循环任务(recurring) | 一次性任务(one-time) |
|---|---|---|
schedule |
{ interval: "1m" } |
不设置或 null |
| 执行完成后 | 更新 runAt 继续循环 |
文档删除或标记 completed |
| 典型用途 | Alerting Rule 检查 | Reporting job、ML 数据同步 |
| Task Manager 行为 | 永不自动删除文档 | 执行完毕后清理文档 |
Alerting Rule 每个 Rule 对应一条循环任务文档,Rule 删除时 Task Manager 同步删除任务文档。Reporting job 是一次性任务,执行完后文档保留一段时间供状态查询,之后由清理机制删除。
容量与并发控制
1 | |
capacity 不足导致任务执行延迟,体现在 schedule_delay_ms 指标上(任务 runAt 到实际开始执行的时间差)。持续高延迟表明 Task Manager worker 数量需要增加,或某类任务占用过多资源。
Task Manager 健康 API
1 | |
响应结构(8.x):
1 | |
drift.p99 超过 poll_interval 的 3-5 倍时,Task Manager 处于压力状态,告警执行会出现肉眼可见的延迟。
实验:直查 .kibana_task_manager 并观察健康状态
第一步:查看所有活跃任务
1 | |
观察 task.status 分布:idle(等待执行)、claiming(正在被 claim)、running(执行中)、failed(失败待重试)。
第二步:查看 Alerting 任务的调度间隔分布
1 | |
第三步:调用健康 API 并记录 drift
1 | |
记录 stats.runtime.value.drift.p99。在 Kibana 空闲状态下该值应低于 poll_interval(默认 3000ms)。
第四步:模拟压力,观察 drift 变化
在 Kibana 中批量创建 50 个间隔 1 分钟的 .es-query Rule,再次调用健康 API,对比 drift.p99 和 load.p99 的变化。
映射到内部对象
| 实验观测 | 内部对象 |
|---|---|
| 任务文档 | .kibana_task_manager 中 type=task 的文档 |
| claim 竞争 | ES update 的 _seq_no 乐观锁 |
| Alerting Rule 周期执行 | task.schedule.interval 驱动的循环任务 |
| Reporting job | task.schedule 为空的一次性任务 |
| 调度延迟 | task.runAt 与实际 startedAt 的差值(drift) |
| 并发上限 | xpack.task_manager.max_workers 配置 |
模式提炼
Task Manager 将调度状态外化到 Elasticsearch,使得 Kibana 本身无状态化:任何实例可以 claim 并执行任何任务,实例宕机后任务会由其他实例在下次轮询时重新 claim(retryAt 机制保证不永久卡在 running 状态)。
这种设计的代价是调度精度受 Elasticsearch 写入延迟和轮询间隔影响,无法达到毫秒级精度。对于告警和报表这类秒级到分钟级的任务来说,这个代价是可接受的。
工程迁移表
| 运维场景 | 处理方法 |
|---|---|
| Task Manager 告警任务延迟高 | 检查 _health API 的 drift;增加 max_workers;减少低价值 Rule 数量 |
| 某任务反复失败 | 查询 task.attempts > 3 的文档;查看 task.state 中的错误信息 |
| 多实例部署任务分配不均 | Task Manager 自动均衡,无需手动干预;检查各实例的 load.p50 是否相近 |
| 清理大量过期 Alerting 任务 | 删除 Rule 时 Task Manager 自动删除对应任务文档;批量删除 Rule 可通过 API |
| 查看特定 Rule 的任务文档 | 任务 _id 格式为 task:alerting:<rule_id>,可直接 GET |
常见误解
Task Manager 的 max_workers 是每个 Kibana 实例的并发上限,不是集群级别的总上限。4 个实例各设 10 workers,集群实际总并发是 40。过高的 max_workers 会导致 Elasticsearch 写入压力增大,因为 claim 操作本身也是写入。
schedule_delay_ms 高不一定是 Task Manager 问题,也可能是 Elasticsearch 响应慢导致轮询 RT 增高。需要结合 Elasticsearch 的 indexing latency 指标一起判断。
任务 status=failed 后不会无限重试,默认最大重试次数由 task type 注册时指定(Alerting 默认 3 次),达到上限后任务进入 dead 状态,需要手动清理或重新启用 Rule。
练习
- 在
.kibana_task_manager中找到某个 Alerting Rule 对应的任务文档,记录_seq_no,然后手动触发一次 Rule 执行,观察_seq_no是否更新。 - 将
xpack.task_manager.poll_interval调整为 10000ms,观察_healthAPI 中configuration.value.poll_interval的变化,以及告警响应延迟的变化。 - 创建一个 Reporting job,在
.kibana_task_manager中确认该任务的task.schedule字段为空(一次性任务的标志)。 - 使用 Kibana Stack Monitoring 中的 Task Manager 面板,对比图表中的 drift 趋势与手动调用
_healthAPI 的返回值是否一致。
系列导航
- 上一篇:深入 Kibana(十一):Reporting——从 Dashboard 到 PDF/PNG 的生成链路
- 下一篇:深入 Kibana(十三):Spaces 与安全——RBAC、Feature Controls 与多租户
参考资料
- Elastic 官方文档,Task Manager - https://www.elastic.co/guide/en/kibana/current/task-manager-production-considerations.html
- Elastic 官方文档,Task Manager health monitoring - https://www.elastic.co/guide/en/kibana/current/task-manager-health-monitoring.html
- Elastic 博客,“Kibana Task Manager: Running Background Tasks at Scale” - https://www.elastic.co/blog/kibana-task-manager-running-background-tasks-at-scale
- Kibana 源码,
x-pack/plugins/task_manager/server/- https://github.com/elastic/kibana
