上一篇展示了 Reporting 作为 Task Manager 的具体用户。这一篇进入 Task Manager 本身,核心问题是:告警、上报、后台任务共用的调度器;任务如何在多 Kibana 实例间分配;与 cron 的区别。

Task Manager 的定位

Task Manager 是 Kibana 进程内的分布式任务调度器,7.4 版本引入,解决的问题是:Kibana 需要运行定期或一次性后台任务(告警检查、报表生成、ML 模型同步等),但不依赖外部消息队列或 cron 守护进程。

1
2
3
4
5
6
7
外部 cron                     Kibana Task Manager
───────────────────────── ──────────────────────────────────
运行在 OS 层面 运行在 Kibana 进程内
任务状态存在本机 任务状态存在 Elasticsearch
多实例需额外协调 多实例通过 ES 乐观锁自动协调
不感知 Kibana 容量 感知 worker 容量,不超载
任务失败需外部监控 内置健康 API 和 drift 监控

所有需要调度的 Kibana 功能——Alerting、Reporting、ML、Fleet——都是 Task Manager 的消费者,注册自己的 task type,由 Task Manager 统一调度。

任务存储:.kibana_task_manager 索引

每个任务对应索引中的一条文档,关键字段如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
{
"_id": "task:abc-123",
"_source": {
"type": "task",
"task": {
"taskType": "alerting:.es-query",
"runAt": "2026-08-06T12:00:00.000Z",
"schedule": { "interval": "1m" },
"params": "{ \"alertId\": \"rule-xyz\" }",
"state": "{ \"previousStartedAt\": \"...\" }",
"ownerId": null,
"status": "idle",
"attempts": 0,
"startedAt": null,
"retryAt": null
}
}
}

ownerId 是 claim 机制的核心:null 表示任务空闲可被抢占,非 null 表示被某个 Kibana 实例持有。

claim 机制与乐观锁

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
Kibana 实例轮询(默认间隔 3s)


查询 .kibana_task_manager
条件:status=idle AND runAt ≤ now AND taskType in registeredTypes
min(capacity_remaining, maxWorkers_per_type) 条


对每条文档执行 ES update:
IF _seq_no == current_seq_no
THEN set ownerId = this_instance_id
set status = claiming
set startedAt = now
ELSE 放弃(另一实例抢先)


claim 成功 → 执行 task handler
执行完成 → 更新 runAt = now + interval,清空 ownerId
执行失败 → 递增 attempts,设置 retryAt(指数退避)

乐观锁通过 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
2
3
4
5
6
7
8
9
10
xpack.task_manager.max_workers (默认 10)

└── 全局并发上限,所有 task type 共享

├── 每个 task type 可设置 maxConcurrency
│ 例如 reporting:printablePdfV2 默认 concurrency=2
│ (Chromium 进程昂贵)

└── 容量不足时任务留在 idle 状态
等待 worker 释放

capacity 不足导致任务执行延迟,体现在 schedule_delay_ms 指标上(任务 runAt 到实际开始执行的时间差)。持续高延迟表明 Task Manager worker 数量需要增加,或某类任务占用过多资源。

Task Manager 健康 API

1
GET kbn:/api/task_manager/_health

响应结构(8.x):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
{
"id": "instance-a",
"timestamp": "...",
"status": "OK",
"stats": {
"configuration": {
"value": { "poll_interval": 3000, "max_workers": 10 }
},
"runtime": {
"value": {
"drift": { "p50": 120, "p99": 450 },
"load": { "p50": 0.2, "p99": 0.8 },
"execution": {
"duration": { "p50": 500, "p99": 3200 }
}
}
},
"workload": {
"value": {
"count": 42,
"task_types": {
"alerting:.es-query": { "count": 30, "status": { "idle": 29, "running": 1 } },
"report:printablePdfV2": { "count": 2, "status": { "idle": 2 } }
}
}
}
}
}

drift.p99 超过 poll_interval 的 3-5 倍时,Task Manager 处于压力状态,告警执行会出现肉眼可见的延迟。

实验:直查 .kibana_task_manager 并观察健康状态

第一步:查看所有活跃任务

1
2
3
4
5
6
7
GET .kibana_task_manager/_search
{
"size": 20,
"query": { "term": { "type": "task" } },
"_source": ["task.taskType", "task.status", "task.runAt",
"task.ownerId", "task.schedule", "task.attempts"]
}

观察 task.status 分布:idle(等待执行)、claiming(正在被 claim)、running(执行中)、failed(失败待重试)。

第二步:查看 Alerting 任务的调度间隔分布

1
2
3
4
5
6
7
8
9
GET .kibana_task_manager/_search
{
"query": { "prefix": { "task.taskType": "alerting:" } },
"aggs": {
"by_interval": {
"terms": { "field": "task.schedule.interval.keyword" }
}
}
}

第三步:调用健康 API 并记录 drift

1
GET kbn:/api/task_manager/_health

记录 stats.runtime.value.drift.p99。在 Kibana 空闲状态下该值应低于 poll_interval(默认 3000ms)。

第四步:模拟压力,观察 drift 变化

在 Kibana 中批量创建 50 个间隔 1 分钟的 .es-query Rule,再次调用健康 API,对比 drift.p99load.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。

练习

  1. .kibana_task_manager 中找到某个 Alerting Rule 对应的任务文档,记录 _seq_no,然后手动触发一次 Rule 执行,观察 _seq_no 是否更新。
  2. xpack.task_manager.poll_interval 调整为 10000ms,观察 _health API 中 configuration.value.poll_interval 的变化,以及告警响应延迟的变化。
  3. 创建一个 Reporting job,在 .kibana_task_manager 中确认该任务的 task.schedule 字段为空(一次性任务的标志)。
  4. 使用 Kibana Stack Monitoring 中的 Task Manager 面板,对比图表中的 drift 趋势与手动调用 _health API 的返回值是否一致。

系列导航

参考资料

  1. Elastic 官方文档,Task Manager - https://www.elastic.co/guide/en/kibana/current/task-manager-production-considerations.html
  2. Elastic 官方文档,Task Manager health monitoring - https://www.elastic.co/guide/en/kibana/current/task-manager-health-monitoring.html
  3. Elastic 博客,“Kibana Task Manager: Running Background Tasks at Scale” - https://www.elastic.co/blog/kibana-task-manager-running-background-tasks-at-scale
  4. Kibana 源码,x-pack/plugins/task_manager/server/ - https://github.com/elastic/kibana