心跳还在,任务却不动了:如何识别分布式系统中的“僵尸 Owner”
在分布式任务协调中,我们经常会使用“租约 + 心跳”机制判断某个任务的执行者是否仍然存活。
一个节点抢到任务后成为 Owner,并定期向 Redis 续租。只要心跳持续更新,其他节点就认为 Owner 仍然健康,不会接管任务。
这个设计看起来没有问题,但它隐藏着一个很容易被忽略的漏洞:
心跳线程还活着,不代表真正执行业务的线程还在推进。
Owner 可能没有宕机,JVM 也仍然正常运行,心跳任务甚至还在每隔几秒续约,但真正的业务线程可能早已卡在数据库查询、线程池、网络请求或 AI 接口中。
这时系统会出现一种特殊故障:
Owner 一直占有执行权,却永远无法完成任务,其他节点也永远没有机会接管。
这就是本文要讨论的 僵尸 Owner 问题。
一句话总结
传统心跳只能证明 Owner 的进程还活着,不能证明任务仍在推进。
更合理的续租条件应该是:
Owner 仍然拥有租约
+
Owner Token 仍然有效
+
业务最近仍有实际进展
说得更直白一点:
不是“活着就续租”,而是“活着并且还在干活,才续租”。
一、普通心跳机制是怎么工作的
假设一个分布式任务同时被多个节点发现:
Node A
Node B
Node C
三个节点竞争执行权,Node A 成功成为 Owner。
Node A:Owner
Node B:Follower
Node C:Follower
系统给 Node A 一段有限时间的租约,例如 30 秒。
为了避免任务执行时间超过 30 秒,Node A 会启动一个心跳线程,每隔 5 秒更新一次租约。
租约时间:30 秒
心跳间隔:5 秒
只要 Node A 还在不断续租,其他节点就认为:
Owner 还活着
→ 任务仍然有人处理
→ 不允许其他节点接管
如果 Node A 宕机,心跳自然停止。
当租约过期后,Node B 或 Node C 就可以重新竞争执行权。
这种设计能够很好地处理:
- 服务器宕机;
- JVM 崩溃;
- 容器被终止;
- 进程被杀死;
- 节点与 Redis 完全断开。
但它处理不了另一类故障:
JVM 没有崩,只有业务线程卡住了。
二、问题到底出在哪里
一个 Owner 内部通常不只有一条线程。
至少可以抽象为两类:
业务线程
负责查询数据、处理参数、调用下游和保存结果
心跳线程
负责定期向共享协调组件续租
正常情况下,两条线程都在工作:
业务线程:不断推进任务
心跳线程:不断续租
发生异常时,可能变成:
业务线程:卡死
心跳线程:正常运行
业务线程可能卡在:
- 数据库查询;
- 分布式锁等待;
- 线程池任务排队;
- HTTP 请求;
- AI 推理接口;
- 文件读写;
- 结果持久化;
- 某个没有超时的阻塞调用。
但只要 JVM 没崩,定时心跳线程仍可能正常运行。
于是 Redis 看到的状态是:
heartbeatAt 一直更新
leaseExpireAt 一直延长
Owner 看起来非常健康
Follower 看到租约一直有效,只能继续等待:
Node B:Owner 仍在续租,不能接管
Node C:Owner 还活着,继续等待
实际上,任务已经几分钟没有任何进展。
最终形成:
Owner 活着,但业务已经停滞;它不能完成任务,却一直阻止其他节点接管。
三、进程存活不等于业务活性
这里需要区分两个概念。
进程存活
表示:
- JVM 还在运行;
- 定时线程还在调度;
- Redis 连接可能正常;
- 心跳仍然可以发送。
通常可以通过:
heartbeatAt
判断。
业务活性
表示:
- 任务正在向完成方向推进;
- 业务阶段仍在变化;
- 已经完成新的子步骤;
- 下游仍然返回有效进度;
- 最近产生了有意义的处理结果。
可以通过:
lastProgressAt
判断。
两者的区别可以总结为:
| 字段 | 证明什么 |
|---|---|
heartbeatAt |
Owner 的进程和心跳机制还在运行 |
lastProgressAt |
Owner 最近确实完成了新的业务进展 |
原来的设计实际上只判断了:
Owner 健康
= 心跳还在
更完整的判断应该是:
Owner 健康
= 进程仍然存活
+ 业务仍在合理时间内推进
四、引入 lastProgressAt
解决僵尸 Owner 的核心思路,是增加一个业务进度时间:
lastProgressAt
它记录:
Owner 最近一次产生有效业务进展的时间。
假设一个任务包含以下流程:
接收任务
↓
查询业务数据
↓
组装请求参数
↓
调用 AI
↓
解析 AI 结果
↓
保存最终结果
可以在关键步骤完成后刷新进度:
数据查询完成
→ 更新 lastProgressAt
参数组装完成
→ 更新 lastProgressAt
收到 AI 响应
→ 更新 lastProgressAt
结果解析完成
→ 更新 lastProgressAt
数据保存完成
→ 更新 lastProgressAt
心跳线程每次准备续租时,不再直接续约,而是先检查:
距离上一次实际业务进展过去了多久?
判断公式可以写成:
now - lastProgressAt <= maxNoProgress
其中:
maxNoProgress
表示当前业务阶段允许的最长无进展时间。
如果超过这个时间:
now - lastProgressAt > maxNoProgress
就说明 Owner 虽然仍然存活,但业务已经长时间没有推进。
此时应该停止续租。
五、停止续租后会发生什么
需要注意,发现 Owner 停滞后,不一定要立刻强制删除它的状态。
更稳妥的过程是:
发现业务长时间无进展
↓
当前 Owner 停止续租
↓
现有租约自然过期
↓
Follower 重新竞争执行权
↓
新的节点成为 Owner
例如:
Node A 长时间无进展
→ Node A 停止续租
租约过期
→ Node B 和 Node C 重新竞争
Node B 成为新 Owner
→ 重新执行任务
为什么不由旧 Owner 直接指定某个 Follower?
因为旧 Owner 当前本身就可能处于异常状态,而且多个节点之间必须通过原子协调机制决定新的 Owner。
如果每个节点自行判断“应该由谁接管”,很容易出现多个新 Owner。
所以正确方式仍然是:
停止续租,让租约过期,然后由候选节点重新竞争。
六、用时间线看得更清楚
假设配置如下:
租约时间:30 秒
心跳间隔:5 秒
最大无进展时间:20 秒
正常情况
00 秒:Node A 成为 Owner
05 秒:完成数据查询,刷新 lastProgressAt
10 秒:完成参数组装,刷新 lastProgressAt
15 秒:心跳检查,最近仍有进展,继续续租
20 秒:收到 AI 结果,刷新 lastProgressAt
25 秒:保存结果完成
业务不断推进,心跳可以正常续租。
业务卡死情况
00 秒:Node A 成为 Owner
05 秒:完成数据查询,刷新 lastProgressAt
06 秒:进入某个阻塞调用
10 秒:距离上次进展 5 秒,允许续租
15 秒:距离上次进展 10 秒,允许续租
20 秒:距离上次进展 15 秒,允许续租
25 秒:距离上次进展 20 秒,到达临界值
30 秒:距离上次进展 25 秒,拒绝续租
随后租约自然过期,其他节点开始重新竞争。
这套机制可以称为:
业务进度感知型租约。
七、什么才算“业务有进展”
引入 lastProgressAt 并不意味着随便刷新时间就可以。
有效进展必须表示:
当前任务距离完成确实更近了一步。
例如:
- 成功读取了一批新数据;
- 完成了一个业务阶段;
- 处理完成一个数据分片;
- 收到新的流式响应;
- 成功写入一个结果批次;
- 完成一个新的子任务;
- 业务状态从一个阶段进入下一个阶段。
以下行为不应该算作有效进展:
- 心跳线程还在执行;
- 日志还在不断输出;
- CPU 仍然在消耗;
- 不断重复同一个失败操作;
- 循环查询同一个状态;
- 不断刷新同一个时间戳;
- 线程还没有退出。
例如下面这种写法毫无意义:
while (true) {
refreshProgress();
}
虽然 lastProgressAt 一直更新,但业务没有真正推进。
这只是把普通 Heartbeat 换了一个名字。
八、最好同时记录当前业务阶段
只记录 lastProgressAt 可以判断任务停滞了多久,但无法快速知道它卡在哪里。
更完整的设计可以同时记录:
currentStage
lastProgressAt
progressVersion
例如:
{
"ownerId": "node-a",
"ownerToken": 12,
"currentStage": "BUILDING_PROMPT",
"progressVersion": 3,
"heartbeatAt": 1710000010000,
"lastProgressAt": 1710000008000
}
其中:
currentStage:当前执行到哪个阶段;lastProgressAt:最近一次进展时间;progressVersion:每次有效进展时递增;ownerToken:当前 Owner 的防护令牌。
如果系统发现:
currentStage = LOADING_CONTEXT
并且 60 秒没有变化
就可以初步判断任务可能卡在数据加载阶段。
这对排查问题非常有价值。
九、不同阶段不能使用同一个超时时间
一个任务中的不同阶段,正常耗时可能差别非常大。
例如:
| 阶段 | 正常耗时 |
|---|---|
| 参数校验 | 1~10 毫秒 |
| Redis 查询 | 5~50 毫秒 |
| 数据库查询 | 20~500 毫秒 |
| Prompt 组装 | 10~200 毫秒 |
| AI 推理 | 3~60 秒 |
| 长报告生成 | 30~180 秒 |
| 保存结果 | 10~1000 毫秒 |
如果所有阶段统一设置:
maxNoProgress = 30 秒
会出现两个问题。
对于参数校验来说,30 秒太长。任务卡死后,要过很久才能发现。
对于大型 AI 报告生成来说,30 秒又太短。正常任务可能被错误地判断为卡死。
因此,更合理的是按阶段配置:
maxNoProgress(currentStage)
例如:
progress-timeout:
validate-request: 1s
load-data: 5s
build-prompt: 2s
call-ai: 90s
parse-result: 5s
persist-result: 10s
心跳续租时,根据当前阶段选择对应的最大无进展时间。
十、超时时间应该如何确定
不能简单地认为:
正常耗时 10 毫秒
→ 配置 30 毫秒
真实生产环境存在很多正常抖动:
- JVM GC;
- 线程调度延迟;
- 网络抖动;
- 数据库锁等待;
- Redis 主从切换;
- 节点 CPU 突然升高;
- 下游短暂排队。
如果阈值设置得太紧,正常但稍慢的任务会被误判为卡死。
更合理的方式是根据历史监控数据确定:
P95
P99
P99.9
例如某阶段的执行耗时:
P50:10ms
P95:25ms
P99:80ms
偶发 GC:200ms
如果只根据平均值设置 30ms,就会误判大量正常请求。
可以采用类似思路:
maxNoProgress
= P99 耗时
+ 网络抖动余量
+ GC 余量
+ 调度延迟余量
或者:
maxNoProgress
= max(固定下限, P99 × 安全系数)
具体数值需要根据实际业务监控不断调整。
十一、调用非流式 AI 时为什么很难判断进度
对于应用内部的前置流程,我们可以在每个关键节点刷新业务进度。
但进入非流式 AI 调用后,问题会变得更加困难。
调用过程是:
应用发送请求
↓
AI 服务内部执行
↓
应用等待完整结果
在最终结果返回前,应用通常只知道:
- 请求已经发出;
- 连接可能仍然存在;
- 结果尚未返回。
应用无法知道模型内部:
- 是否正在排队;
- 是否已经开始推理;
- 已经完成多少;
- 是否卡在供应商内部;
- 是否即将返回;
- 是否已经发生内部错误。
对于这种纯阻塞黑盒调用,应用层无法准确判断“它是否仍在推进”。
能做的主要是:
- 设置连接超时;
- 设置读取超时;
- 设置请求总超时;
- 根据供应商 SLA 设置阶段期限;
- 使用异步任务接口;
- 查询下游任务状态;
- 超时后取消请求;
- 必要时触发重新接管。
所以在非流式 AI 阶段:
lastProgressAt不能提供连续进度,只能依靠阶段级硬超时控制。
十二、流式 AI 为什么更容易判断进度
如果使用流式 AI 接口,服务端会不断返回数据:
chunk 1
chunk 2
chunk 3
chunk 4
每收到一个有效 Chunk,就可以刷新:
lastProgressAt
例如:
10:00:00 收到第一个 Token
10:00:01 收到新的一批 Token
10:00:02 再次收到新内容
这说明下游确实还在推进。
如果长时间没有新 Chunk:
now - lastChunkAt > streamIdleTimeout
就可以认为:
- 流可能中断;
- 网络可能阻塞;
- 下游可能卡死;
- 当前请求需要取消或重新处理。
流式调用通常还需要区分三个超时:
首 Token 超时
流空闲超时
总执行超时
首 Token 超时
请求发出后,最长允许多久收到第一个 Token。
流空闲超时
已经开始返回数据后,两次 Chunk 之间最长允许间隔多久。
总执行超时
整个流式任务最多可以运行多久。
这样比单独设置一个统一超时更加合理。
十三、为什么还需要 Fencing Token
系统停止给旧 Owner 续租,并不代表旧 Owner 的业务线程一定停止。
可能出现下面的情况:
Node A 长时间没有进展
→ 系统停止给它续租
→ 租约过期
→ Node B 接管任务
→ Node A 突然恢复
现在 Node A 和 Node B 都可能继续执行。
因此,每次 Owner 重新选举时,都需要生成一个更大的版本号:
Node A:ownerToken = 10
Node B:ownerToken = 11
最终写入结果时必须校验:
提交 Token = 11
当前有效 Token = 11
→ 允许写入
提交 Token = 10
当前有效 Token = 11
→ 拒绝写入
这个递增版本号就是:
Fencing Token
它的作用不是让旧 Owner 立刻停止执行,而是:
即使旧 Owner 后来恢复,也不能再覆盖新 Owner 的最终结果。
所以系统真正保证的是:
旧 Owner 可以继续运行
但旧 Owner 不能提交过期结果
十四、这套方案能保证绝对只执行一次吗
不能。
假设旧 Owner 已经调用了外部 AI,随后被判断为卡死。
新 Owner 接管后,也可能再次调用 AI。
这时外部服务可能收到两次请求。
Fencing Token 只能保护受我们控制的最终写入,例如:
- 数据库结果;
- Redis 状态;
- 任务完成标记。
它不能撤销已经发生的外部调用。
因此,这套机制能保证的是:
- 避免僵尸 Owner 永久占用执行权;
- 允许任务在故障后恢复;
- 阻止旧 Owner 脏写;
- 尽量减少重复执行。
它不能保证:
任何故障情况下都绝对只调用一次
对于具有副作用的下游操作,还需要配合:
- 幂等 Key;
- 唯一请求 ID;
- 数据库唯一约束;
- 下游状态查询;
- 重复调用检测;
- 补偿机制。
十五、自动接管也可能误判
判断 Owner 长时间没有进展后,允许其他节点接管,看起来很合理。
但旧 Owner 不一定真的卡死,也可能只是暂时变慢。
例如:
- Full GC;
- 数据库短暂抖动;
- 网络延迟;
- AI 服务排队;
- 操作系统调度延迟;
- 下游限流。
如果 maxNoProgress 设置得太小:
正常慢任务
→ 被误判为僵尸 Owner
→ 新节点接管
→ 两个节点重复执行
如果设置得太大:
真正卡死的任务
→ 很长时间后才能恢复
这是一个典型的权衡:
阈值太短
→ 恢复快,但误判多
阈值太长
→ 误判少,但恢复慢
可以采用以下方式降低误判:
- 连续多次检测无进展后才停止续租;
- 根据阶段设置不同超时;
- 设置最大任务总时长;
- 接管前增加短暂宽限期;
- 使用 P99 或 P99.9 历史耗时;
- 对高风险任务只告警,不自动接管;
- 对低风险、高成本任务允许自动接管。
十六、推荐的 Owner 健康判断
可以把 Owner 状态分成四类:
| 心跳状态 | 业务进度 | 判断 |
|---|---|---|
| 正常 | 正常 | Owner 健康,继续续租 |
| 异常 | 未知 | 节点可能失效,等待租约过期 |
| 正常 | 超时 | 僵尸 Owner,停止续租 |
| 异常 | 超时 | Owner 明显失效,允许后续接管 |
从概念上,可以把续租条件表达为:
leaseStillOwned
AND tokenStillCurrent
AND progressWithinDeadline
如果还需要单独检查心跳状态,可以写成:
heartbeatHealthy
AND progressWithinDeadline
不过实际实现中,心跳线程本次能够正常执行续租检查,本身已经说明当前心跳线程仍然存活。
真正新增的关键判断是:
progressWithinDeadline
也就是:
now - lastProgressAt
<=
当前阶段允许的最大无进展时间
十七、一个更完整的状态设计
共享状态可以设计成:
{
"status": "RUNNING",
"ownerId": "node-a",
"ownerToken": 12,
"currentStage": "CALLING_AI",
"progressVersion": 4,
"heartbeatAt": 1710000010000,
"lastProgressAt": 1710000008000,
"leaseExpireAt": 1710000040000
}
字段含义如下:
| 字段 | 作用 |
|---|---|
status |
当前任务状态 |
ownerId |
当前执行节点 |
ownerToken |
当前 Owner 的 Fencing Token |
currentStage |
当前业务阶段 |
progressVersion |
有效进展次数或版本 |
heartbeatAt |
最近心跳时间 |
lastProgressAt |
最近业务推进时间 |
leaseExpireAt |
当前租约过期时间 |
心跳续租时需要校验:
任务仍然属于当前 Owner
Token 仍然是最新版本
当前任务仍然处于 RUNNING
业务最近仍有进展
全部满足后才允许延长租约。
十八、通用执行流程
可以把整个机制整理成下面的流程:
节点成为 Owner
↓
开始执行任务
↓
完成关键业务阶段
↓
更新 currentStage
更新 progressVersion
更新 lastProgressAt
↓
心跳线程准备续租
↓
检查 Token 和租约归属
↓
检查最近业务进度
├── 仍在合理时间内
│ ↓
│ 继续续租
│
└── 长时间无进展
↓
停止续租
↓
租约自然过期
↓
其他节点重新竞争
↓
新 Owner 获得更高 Token
旧 Owner 后续即使恢复,也会因为 Token 过期而无法写入最终结果。
十九、这套设计适合哪些任务
业务进度感知型租约适合:
- AI 推理;
- 报表生成;
- 文件转换;
- 视频处理;
- 大规模数据计算;
- 搜索索引构建;
- 分布式爬取;
- 批处理任务;
- 工作流执行;
- 长时间第三方接口调用。
这些任务通常具备以下特点:
执行时间长
+
中途可能卡住
+
需要故障接管
+
不希望旧节点脏写
对于执行时间只有几毫秒的简单请求,加入完整的阶段进度和租约机制可能没有必要。
二十、需要监控哪些指标
上线后建议监控:
owner_heartbeat_total
owner_lease_renew_total
owner_lease_renew_rejected_total
owner_progress_update_total
owner_no_progress_timeout_total
owner_takeover_total
owner_fencing_rejected_total
owner_stage_duration
owner_task_total_duration
owner_ai_first_token_duration
owner_stream_idle_timeout_total
重点观察:
- 哪些阶段最容易无进展;
- 无进展接管是否频繁;
- 是否存在大量错误接管;
- 某个阶段的 P99 是否持续上升;
- Fencing 拒绝是否经常发生;
- AI 首 Token 是否越来越慢;
- 任务平均恢复时间是多少。
只有通过监控数据,才能合理调整 maxNoProgress。
总结
传统 Heartbeat 解决的是:
Owner 的进程是否还活着
业务进度检测解决的是:
Owner 是否仍然在真正推进任务
当业务线程卡死、但心跳线程正常时,如果继续无条件续租,就会产生一个永远占有执行权的僵尸 Owner。
解决思路是引入:
lastProgressAt
currentStage
progressVersion
maxNoProgress
Fencing Token
将续租逻辑从:
只要心跳还在
→ 继续续租
升级为:
Owner 仍然拥有租约
并且 Token 仍然有效
并且业务最近仍有实际进展
→ 才允许继续续租
如果业务长时间没有推进:
停止续租
→ 等待租约过期
→ 其他节点重新竞争
→ 新 Owner 使用更高 Token
最终可以用一句非常直白的话概括这套设计:
Heartbeat 只能证明 Owner 还活着,
lastProgressAt才能证明它还在干活。