5842 字
约 19 分钟
1
如何识别分布式系统中的“僵尸 Owner”

心跳还在,任务却不动了:如何识别分布式系统中的“僵尸 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 才能证明它还在干活。

如何识别分布式系统中的“僵尸 Owner”
http://www.clxhxhhr.top/posts/129/
作者
clxstart
发布于
2026-07-17
许可协议
CC BY-NC-SA 4.0
评论
0 条
还没有评论,先写一条吧。