跳至正文
来两杯美式
返回

给生产 Agent 做体检(五):不要消息队列——基于数据库的任务编排与可靠性

By 来两杯美式
发布于

全量对话诊断是一个典型的长时任务:一次任务要跑完前几篇讲过的质检、风险筛选、结构化诊断、归一化修复整个流水线,中间穿插几十次 LLM 调用,跑上几十分钟是常态。把这种任务放进 Web 请求里同步等待显然不现实,于是很自然地会想到”该上消息队列了”。

但这一篇的结论恰恰相反:在中等吞吐量级下,已有的 PostgreSQL 加上事务内的 compare-and-set 原子领取,加上心跳与租约恢复,就是一个更简单也更可靠的编排层。不需要引入 Kafka 或 RabbitMQ,不需要独立部署和监控一个 broker,本地开发一条命令就能起,生产多实例部署时天然支持并发领取和故障接管。把”队列”做成数据库里的一张任务表,可靠性语义反而更清晰——任务状态、归属、进度全在一个事务边界里,不用操心消息与数据库的双写一致性。

当然这个选择有明确的适用边界。本文会先拆解这套编排的设计细节,再专门用一节讨论”为什么不用 MQ”,把什么时候必须上 MQ 的边界说清楚,避免这套方案被误用到它撑不住的场景。

数据库任务编排总览:中心是 PostgreSQL 任务表与状态机,多 Pod 通过 CAS 原子领取,配合心跳租约接管、内嵌 Worker、受控并发与任务归属安全

一、需求:长时任务、多实例、本地一条命令

先把问题定义清楚。全量诊断任务有三个刚性约束:

第三条看起来只是开发体验问题,但它深刻影响架构选型。如果编排层依赖一个独立的消息中间件,本地就得拉起这个中间件;如果 worker 是独立进程,本地就得跑两个进程。每多一个必须存活的组件,开发和部署的失败模式就多一类。而数据库是无论如何都绕不开的——任务状态总得持久化在某处。既然任务表已经在数据库里了,“队列”无非是这张表上的一个领取协议。

二、数据库 CAS 领取:用一条 UPDATE 实现原子分配

核心机制非常朴素:任务表里有一列 status,worker 领任务就是一条条件更新。

UPDATE optimizer_analysis_task
SET status = 'running',
    claimed_at = now(),
    heartbeat_at = now()
WHERE id = (
  SELECT id FROM optimizer_analysis_task
  WHERE status = 'pending'
  ORDER BY created_at
  FOR UPDATE SKIP LOCKED
  LIMIT 1
)
RETURNING id, payload;

这条语句在事务内完成 compare-and-set:只有 status 仍然是 pending 时更新才会生效,RETURNING 把领到的任务返回给执行方。PostgreSQL 的 FOR UPDATE SKIP LOCKED 让并发的领取者跳过已被锁住的行,不互相阻塞。

多 Pod 下的时序是这样的:

  1. 两个 Pod 的 worker 同时轮询,都看到队列里有一条 pending 任务。
  2. 两个 Pod 各自尝试原子更新同一条任务的 status
  3. 数据库行锁保证只有一个 UPDATE 成功,成功的 Pod 拿到 RETURNING 结果并开始执行;失败的 Pod 得到零行更新,继续轮询。
  4. 如果队列里有多条不同任务,两个 Pod 可以各自领到不同的任务并发执行,每个 Pod 同时最多处理一个任务。

这个协议的美妙之处在于它没有引入任何新的协调组件

任务状态机也全部落在一张表里,pending → running → 终态,每个状态迁移都是一次数据库写入,天然可查询、可审计。这里的 pending / running 对应第四篇 status 词表里的 PENDING / RUNNING;终态同样是那一张表的划分——COMPLETED / PARTIAL / INSUFFICIENT_DATA / FAILED,下文 SQL 里简写为 completed / failed。排查问题时直接查表,比翻 broker 的日志舒服得多。

三、内嵌 Worker:把轮询循环塞进 Web 进程

有了 CAS 领取协议,下一步是决定谁来执行循环。答案是:Web 进程内直接起一个轮询循环

// pseudo-code
async function workerLoop() {
  while (!shutdownRequested) {
    const task = await claimNextAnalysisTask(); // CAS claim
    if (task) {
      await runAnalysisTask(task); // errors caught at loop boundary
    } else {
      await sleep(POLL_MS);
    }
    // next iteration scheduled only after previous settles
  }
}

几个细节值得展开:

串行迭代,单 Pod 永不重叠消费。 循环只有在上一次迭代完全 settle(成功或失败)之后才调度下一次。这意味着单个 Pod 在任意时刻最多处理一个任务。并发度等于 Pod 数,这个模型足够直观:想提高吞吐就加 Pod,每个 Pod 的资源消耗可预期。

立即执行一次迭代。 循环启动时先尝试领取一次,而不是先睡一个轮询间隔。这让本地开发体验更好——创建任务后几乎立刻开始处理,而不是等一个空转周期。

process 级单例防热重载。 本地开发时模块会被反复重新求值,如果每次都起新循环,就会有好几个循环同时抢任务。用一个挂在 globalThis 上的单例标记解决:

const g = globalThis as { __optimizerWorkerStarted?: boolean };
if (g.__optimizerWorkerStarted) return;
g.__optimizerWorkerStarted = true;

运行时守卫。 这个 bootstrap 绝不能在构建期执行、不能跑进浏览器 bundle、不能在 Edge runtime 里启动,也不能污染测试模块。所以入口处有一组环境判断,只在 Node.js server runtime 里生效。这类守卫平时不起眼,一旦缺失,最轻的后果是 build 报错,最重的后果是测试环境里有幽灵 worker 在偷偷领任务。

可用配置关闭。 环境变量 EMBEDDED_WORKER_ENABLED 默认开启,但保留关闭的能力——特殊部署形态如果想要专用的 worker 实例(比如给跑任务的 Pod 单独配资源),可以关掉 Web 进程里的循环,回退到独立 worker 脚本。内嵌是默认路径,不是唯一路径。

这样一来,部署拓扑收敛为:本地 pnpm dev 一条命令;生产每个 Pod 跑同一个应用镜像,内嵌 worker 随进程启动。

四、故障恢复:心跳与租约

长时任务最怕两件事:进程挂了任务卡死,以及优雅关闭时在途请求被腰斩。

心跳。 任务处理过程中会周期性刷新 heartbeat_at,尤其在每个问题簇处理完成时刷新。心跳是”这个任务还活着有人在跑”的证据。

租约接管。 轮询领取的条件不只有 status = 'pending',还包括”running 状态但心跳超过阈值”的任务:

WHERE status = 'pending'
   OR (status = 'running' AND heartbeat_at < now() - interval 'TASK_STALE_MS')

某个 Pod 在任务跑到一半崩溃后,它的任务会停在 running 状态;超过 stale 阈值后,其他 Pod 的轮询会把它当作可领取任务重新 CAS、重新执行。故障恢复不需要人工干预,也不需要额外组件,只是领取条件多了一个分支。 任务本身设计为幂等重跑——重新执行只是重复 LLM 调用和重写结果,不产生脏状态。

优雅关闭。 进程收到关闭信号时,循环停止调度下一轮,但不强杀在途的 LLM 请求——硬掐一个正在生成的请求,除了浪费已经花掉的 token 之外没有任何收益。正在处理的任务随进程一起中断,之后由租约机制接管。换句话说:正常关闭负责”不再领新任务”,异常中断交给租约兜底,两者配合覆盖了所有退出路径。

五、为什么不用消息队列

这是全文最容易被讨论的一节,把取舍写透。

先看我们实际面对的任务量级:全量诊断是一个低频、批量、人工触发的操作。任务不是每秒成千上万条地涌入,而是偶尔来一批,每条还要跑几十分钟。在这个量级下,消息队列的核心卖点——高吞吐管道——完全没有用武之地。瓶颈在 LLM 调用耗时,不在投递速率。

再逐项对比:

维度数据库任务表消息队列
部署组件复用已有 PostgreSQL新增 broker 及其运维
本地启动一条命令需要拉起中间件
与业务数据一致性同库同事务,无双写消息与 DB 双写,需处理不一致
任务状态查询直接 SQL 查表需另建状态存储或查 broker 内部
吞吐上限中等(百级/秒以下舒适)极高
延迟轮询间隔级(秒级)毫秒级推送
削峰填谷无(靠并发上限)天然缓冲

其中”一致性”一项最容易被低估。用 MQ 编排时,“任务已写入数据库”和”消息已发出”是两个系统里的两次写入,要么接受最终一致带来的各种边界情况,要么引入事务性发消息的复杂度。而数据库方案里,创建任务和暴露给领取者就是同一个事务提交,不存在中间态。

更要紧的是可靠性语义并没有损失。很多人下意识觉得”MQ 更可靠”,但在任务编排场景里,我们需要的可靠性是:任务不丢、不重复执行到一半没人管、挂了能恢复。这三条 CAS 领取加租约全都覆盖了,而且恢复逻辑就在业务代码里,可测试、可调试。MQ 提供的 at-least-once 投递同样要求消费端幂等,并没有替你解决最难的那部分。

那么什么时候 MQ 真正不可替代?我认为边界在这几条:

诊断任务的规模离每一条都很远。工程上引入一个组件的理由应该是”它解决了我们真实存在的问题”,而不是”这个场景大家一般都这么做”。

六、受控并发:Action 生成阶段

编排层解决了”任务被谁执行”,任务内部还有一个”并发度怎么控制”的问题。诊断流水线中,Action 生成阶段要对多个问题簇分别调用 LLM,串行执行会显著拉长总耗时,因此这一阶段有独立的并发配置。

设计上有几条纪律:

候选选择是纯函数。 先按类型筛出全部 RISKACTION 两组,各自按 priorityScore 降序,保留全部 RISK 加优先级最高的 10 个 ACTION。这里修掉过一个 bug:早期实现用统一的 top-10 截断,导致 RISK 占用了 ACTION 的名额——两类发现混在一个池子里排序,语义就错了。

并发可配置但有硬上限。 环境变量控制并发数,默认 2,上限 4。非法值回退默认值而不是报错拒绝启动。上限与诊断阶段一致,防止误配置打出超出预期的请求压力。注意这只是单任务内的并发限制,不是跨进程全局限流——多 Pod 同时跑多个任务时,总在途请求数是两者相乘,这是已接受的边界。

完成顺序不影响结果。 并发调用完成的先后是随机的,但每个成功结果都带着原候选索引写入内存结果集,全部完成后按原索引排序、依次赋予从 1 开始的连续 rank。网络响应的乱序因此完全不改变报告内容。失败列表也按候选顺序输出。

失败隔离。 单个问题簇的 LLM 错误在并发任务内部捕获,带上下文记录为任务 warning,不抛给并发调度器,也不阻断其他簇。每个簇无论成败都在 finally 里触发心跳刷新。唯一例外的”基础设施错误”是心跳持久化本身失败——这种错误继续向外抛并使任务失败,因为心跳失效意味着租约机制的前提被破坏,静默吞掉比失败更危险。

七、任务归属:把 owner 放进查询条件,而不是路由层

诊断任务包含用户上传的对话数据,归属安全不能只做一层。这一节的原则是:ownerUserId 是数据库查询条件的一部分,而不是路由层检查一次就完事的装饰

具体来说,任务表新增必填列 owner_user_id,外键关联用户表,配 (owner_user_id, created_at) 复合索引支撑按人分页。API 路由从 session 取出用户 ID 后显式传入业务函数,所有面向用户的函数签名都强制要求 owner:

createAnalysisTask(input, ownerUserId);
listAnalysisTasks(options, ownerUserId);
getAnalysisProgress(taskId, ownerUserId); // WHERE id = ? AND owner_user_id = ?
deleteAnalysisTask(taskId, ownerUserId); // zero rows deleted => 404

关键在”查询条件”四个字。如果归属检查是独立的一步——先查任务、比对 owner、再决定是否返回——那么每一个新增的读取路径都要记得补这个检查,漏一处就是越权。而把 owner 直接拼进 WHERE 子句后,“忘记过滤”在结构上就不可能:不传 owner 的查询根本写不出来。这是把安全从”纪律”降级成”类型”的做法。

几个配套细节:

八、观测性:安静地空转,吵闹地报错

内嵌 worker 的日志策略是两条相反的规则:

任务级失败则由任务执行函数持久化到任务记录里,用户在任务详情里可见,不依赖日志。另外,独立 worker 脚本被保留下来,且保留了详细的交互式输出——它不再是部署必需品,但在手工补偿、复现排障、“用一个进程单独跑这个任务看看会发生什么”的场景下,一个话多的 worker 是无价的调试工具。

九、踩坑清单

实现这套编排过程中踩过的坑,按伤害程度排列:

十、局限与适用边界

这套方案不是银弹,边界要自己先说清楚:

十一、总结

这一篇把视线从”单个任务内部怎么跑”拉到了”任务之间怎么被编排”。核心结论是:在中等吞吐下,PostgreSQL 事务加 CAS 领取加租约恢复,用一张任务表就撑起了长时异步任务的全部可靠性需求——原子分配、故障接管、优雅关闭、状态可查——同时把部署拓扑压到最简:一个数据库、一个应用镜像、一条本地启动命令。消息队列是伟大的工具,但它解决的问题(高吞吐、多消费组、削峰、跨语言)在诊断场景里一个都不存在,为不存在的需求引入真实的运维成本,是架构上典型的负收益。

回扣系列主线:模型负责语义判断,应用负责协议边界。任务编排就是”协议边界”在执行层的体现——谁拥有任务、谁可以读它、它挂了由谁接管、并发完成顺序如何不污染结果,这些都不是模型能兜底的,只能靠应用层的机制硬性保证。前面四篇解决了”一次诊断怎么跑对”,这一篇解决了”很多次诊断怎么跑稳”。

下一篇是系列的最后一篇:诊断明细堆在数据库里只是中间产物,怎么把它变成 PM 真正会看、看得懂、能据此决策的报告——从字段取舍到信息层级的产品化过程。


分享这篇文章:
通过邮件分享这篇文章✓ 链接已复制
查看系列全部文章
  1. 01.把存量系统交给 Agent,先想清楚 API、Skill、MCP 各是什么
  2. 02.存量系统 Agent 化,物理上到底怎么「接」
  3. 03.单点跑通之后,怎么把 Agent 能力沉淀成全公司可复用
  4. 04.DeepSeek Harness:把 Agent 宿主变成可组合的基础设施,企业能拿它做什么
  5. 05.一个系统,两种访客:给浏览器和 Agent 设计同一套身份体系
  6. 06.定时是 Agent 平台的另一半:Trigger/Action Registry 与调度原语分层
  7. 07.Agent 刚才干了什么:审计、调用日志与敏感数据豁免
  8. 08.提示词安全攻防全解析:越狱、注入与信息泄露
  9. 09.给生产 Agent 做体检(一):只用对话文本的黑盒诊断工作流
  10. 10.给生产 Agent 做体检(二):诊断之前,先定义人口——数据质检与风险筛选
  11. 11.给生产 Agent 做体检(三):让 LLM 稳定输出结构化诊断——严格 JSON Schema 实践
  12. 12.给生产 Agent 做体检(四):LLM 输出不完美怎么办——归一化、一次修复与失败隔离
  13. 13.给生产 Agent 做体检(五):不要消息队列——基于数据库的任务编排与可靠性
  14. 14.给生产 Agent 做体检(六):从诊断明细到 PM 能用的报告——评估结果的产品化
  15. 15.AI Agent:从工具到同事,中间隔着一层「自主性」
  16. 16.Agent 评测工程化(一):为什么不能只看平均分
  17. 17.Agent 评测工程化(二):Rubric 不是 Prompt,而是可执行的质量协议
  18. 18.Agent 评测工程化(三):从低分样本到问题簇
  19. 19.Agent 评测工程化(四):红线、一票否决与 1-5 分
  20. 20.Agent 评测工程化(五):让 LLM-as-judge 稳定输出结构化结果
  21. 21.Agent 评测工程化(六):评测前,先把 Agent 输出拆开
  22. 22.Agent 评测工程化(七):从评分报表到 PM 决策工作台
  23. 23.Agent 评测工程化(八):不用消息队列,也能跑长任务
  24. 24.Agent 评测工程化(九):把一次性评测变成持续优化体系

上一篇
给生产 Agent 做体检(六):从诊断明细到 PM 能用的报告——评估结果的产品化
下一篇
给生产 Agent 做体检(四):LLM 输出不完美怎么办——归一化、一次修复与失败隔离