跳到正文
Elaine Blog
返回

Worker 怎样调度任务:队列、租约、心跳与背压

更新于:
Agent Runtime

任务已经批准,队列里有一条执行消息。Worker A 领取后卡住,系统让 Worker B 接管;几秒后 A 又恢复。这时两者都可能认为自己应该执行退款。

队列分发、任务所有权和业务幂等分别解决不同问题。把消息交给一个 Worker,不代表它永远是唯一执行者;给任务加租约,也不代表旧进程会立即停止。

消息确认意味着什么

消费者收到消息后,如果立刻确认再执行,执行前崩溃可能丢失任务;执行后再确认,则在执行成功、确认丢失时可能重复收到。因此常见设计容许重复交付,并让业务处理可以识别同一次任务或操作。

队列消息最好引用持久任务,而不是携带唯一的一份完整业务状态。Worker 领取后重新读取任务当前版本,防止旧消息让已取消任务继续执行。

消息存活期、可见性超时和任务截止时间也不同。一条消息暂时对其他消费者不可见,不代表业务授权在这段时间内永久有效。

Lease 解决暂时所有权

租约记录 owner、expires_at 和一个单调递增的 generation。Worker 定期续期,正常完成时释放;失联后租约到期,另一个 Worker 可以获取。

心跳只能表示执行者仍在联系,不能证明它持续取得有效业务进展。一个陷入循环的 Worker 也能正常发心跳,所以还需要任务总时间、步骤预算和进展检查。

租约时间过短容易误判慢执行,过长则增加接管等待。合理设置依赖心跳频率、网络抖动和任务特征,不存在所有系统通用的几秒钟数值。

为什么还需要 fencing token

A 获得 generation 41,B 接管后获得 42。若存储或副作用入口拒绝来自 41 的新写入,A 恢复后就无法继续覆盖当前状态。

A 领取 41 → A 暂停 → 租约到期 → B 领取 42
A 恢复并提交 41 → 下游比较当前代数 → 拒绝旧执行者

只生成数字但没有下游检查,不会产生保护。对没有 fencing 支持的外部支付,仍然需要稳定 operation ID、参数绑定和网关幂等;两种机制不能互相替代。

lease_demo.py下载用确定性时钟展示代数变化及拒绝旧写入,不启动真实多机系统。实际系统还需考虑数据库时间、原子更新和网络分区。

超时应该分层

排队超时限制等待资源的时间;执行超时限制开始后的工作;单次请求超时限制某个依赖;总截止时间约束整个任务。只有单次 HTTP 超时,无法限制无限重试累计数小时。

父任务给子任务分配的截止时间不能超过自身剩余时间。超时后取消协程,也不一定停止线程、子进程或远端服务。Runtime 需要使用具体执行设施的取消接口,并保留未知结果。

Temporal 的 Activity 有不同超时概念,源码实践会解释 Schedule-to-Start、Start-to-Close、Schedule-to-Close 与心跳相关超时。它们不能简单映射成一个通用 timeout 参数。

背压怎样保护下游

假设模型服务每分钟只允许一百个请求,Worker 却每秒领取一百个任务,大量请求只会排在重试里。背压需要在进入昂贵阶段前控制速率、并发和队列容量。

并发限制与速率限制也不相同。十个请求同时执行,可能每秒完成上百个;每分钟一百次请求,也可能在第一秒全部发出。系统可能同时需要两者。

租户配额避免单个用户占满所有执行资源。优先级可以改善紧急任务响应,但长期低优先级任务可能饿死,需要公平或老化策略。拆出长短任务队列也可以减少队首阻塞,但增加运维复杂度。

重试队列与死信队列不是垃圾桶

临时错误按退避重新调度,不必让 Worker 原地长时间睡眠。重试要保留原任务和操作身份,记录尝试次数与最后错误。

死信保存无法自动处理的任务。人工重新投递前应核对外部动作状态,不能把支付结果未知的消息当成从未执行。批量重放还会给下游造成突发压力,需要再次限流。

发布和停机时怎样交接

优雅停机通常先停止领取新任务,再等待或取消在途步骤,保存必要状态并释放资源。进程最终仍可能被强制结束,因此可靠性不能只依赖关闭钩子一定执行。

异步子任务、沙箱和临时凭证需要外部回收机制。Worker 崩溃后,父进程的 finally 不会替远端环境自动清理。

调度记录应关联 task、run、owner、generation 和操作标识。这样排查重复执行时,能够区分消息重复、旧 Worker 写入和业务接口不幂等。下一篇讨论把这些执行事实转换为用户可恢复的进度视图


分享这篇文章:

上一篇
人工介入怎样恢复:澄清、审批、编辑与接管
下一篇
执行过程怎样交给用户:事件流、断线重连与取消