全量对话诊断是一个典型的长时任务:一次任务要跑完前几篇讲过的质检、风险筛选、结构化诊断、归一化修复整个流水线,中间穿插几十次 LLM 调用,跑上几十分钟是常态。把这种任务放进 Web 请求里同步等待显然不现实,于是很自然地会想到”该上消息队列了”。
但这一篇的结论恰恰相反:在中等吞吐量级下,已有的 PostgreSQL 加上事务内的 compare-and-set 原子领取,加上心跳与租约恢复,就是一个更简单也更可靠的编排层。不需要引入 Kafka 或 RabbitMQ,不需要独立部署和监控一个 broker,本地开发一条命令就能起,生产多实例部署时天然支持并发领取和故障接管。把”队列”做成数据库里的一张任务表,可靠性语义反而更清晰——任务状态、归属、进度全在一个事务边界里,不用操心消息与数据库的双写一致性。
当然这个选择有明确的适用边界。本文会先拆解这套编排的设计细节,再专门用一节讨论”为什么不用 MQ”,把什么时候必须上 MQ 的边界说清楚,避免这套方案被误用到它撑不住的场景。

一、需求:长时任务、多实例、本地一条命令
先把问题定义清楚。全量诊断任务有三个刚性约束:
- 长时。单任务耗时以分钟到几十分钟计,远超 HTTP 请求的生命周期,必须有异步执行和状态查询机制。
- 生产多实例部署。线上是多 Pod 部署,任意一个 Pod 都可能承载创建任务的请求,也任意一个 Pod 都应该能执行任务;某个 Pod 挂掉后,它在跑的任务要能被别的 Pod 接管。
- 本地开发要简单。开发者拉下仓库后,希望
pnpm dev一条命令就跑起完整功能,而不是”先起数据库、再起 Redis、再起 worker、最后起 Web”四步走,其中任何一步配错都会让人怀疑代码有 bug。
第三条看起来只是开发体验问题,但它深刻影响架构选型。如果编排层依赖一个独立的消息中间件,本地就得拉起这个中间件;如果 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 下的时序是这样的:
- 两个 Pod 的 worker 同时轮询,都看到队列里有一条
pending任务。 - 两个 Pod 各自尝试原子更新同一条任务的
status。 - 数据库行锁保证只有一个 UPDATE 成功,成功的 Pod 拿到
RETURNING结果并开始执行;失败的 Pod 得到零行更新,继续轮询。 - 如果队列里有多条不同任务,两个 Pod 可以各自领到不同的任务并发执行,每个 Pod 同时最多处理一个任务。
这个协议的美妙之处在于它没有引入任何新的协调组件:
- 不需要 leader election——没有”谁是主 worker”的选举和脑裂问题;
- 不需要独立 worker 进程作为可用性依赖——任何活着的 Pod 都是 worker;
- 水平扩展是自动的——多起一个 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,串行执行会显著拉长总耗时,因此这一阶段有独立的并发配置。
设计上有几条纪律:
候选选择是纯函数。 先按类型筛出全部 RISK 和 ACTION 两组,各自按 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 的查询根本写不出来。这是把安全从”纪律”降级成”类型”的做法。
几个配套细节:
- 越权与不存在统一 404。访问别人的任务和访问不存在的任务返回同样的 404,不返回 403。403 等于向攻击者确认”这个任务 ID 是有效的”,为任务 ID 枚举提供了 oracle。
- 两个 userId 不能混用。上传数据里原始对话行自带一个业务
userId,那是被诊断对话数据的属性,用于分析维度,不承担任何系统权限语义;任务的ownerUserId是系统登录用户,是权限的唯一依据。名字相似,概念完全不同,混用必然出安全事故。 - Worker 不受用户归属约束。Worker 是系统内部调用方,按已认领的
taskId执行流水线,不需要也不应该模拟用户 session。claim 成功即代表系统已接管该任务,owner 不进入 claim 和处理逻辑。权限边界属于 API 面,不属于执行面。
八、观测性:安静地空转,吵闹地报错
内嵌 worker 的日志策略是两条相反的规则:
- 空闲轮询不打日志。轮询间隔是秒级的,每次空转都打一行的话,本地终端和线上日志都会被无意义的”no task”淹没,真正需要看的信号反而被冲掉。
- 意外错误带稳定前缀打到 stderr。循环边界的意外异常统一以固定前缀输出(如
optimizer worker failed:),之后照常进入下一轮等待。前缀稳定意味着可以按它配置告警和检索,Web 服务本身不受影响。
任务级失败则由任务执行函数持久化到任务记录里,用户在任务详情里可见,不依赖日志。另外,独立 worker 脚本被保留下来,且保留了详细的交互式输出——它不再是部署必需品,但在手工补偿、复现排障、“用一个进程单独跑这个任务看看会发生什么”的场景下,一个话多的 worker 是无价的调试工具。
九、踩坑清单
实现这套编排过程中踩过的坑,按伤害程度排列:
- 热重载起多个循环。开发模式模块重载后旧循环未死、新循环又起,多个循环互相抢任务,症状是任务被处理两次、日志交错。process 级单例标记是唯一的解。
- worker bootstrap 跑进了 build 或测试。构建产物里出现了数据库连接尝试,或测试环境有幽灵 worker 领走了测试任务。运行时守卫必须显式白名单”Node.js server runtime”,而不是黑名单排除。
- 统一 top-10 截断混排两类候选。
RISK和ACTION混在一个池子里排序截断,导致高风险项挤掉了可执行项的名额。候选选择拆成纯函数后才有单测兜住。 - 并发完成顺序泄漏进 rank。直接按完成顺序赋 rank,同样的输入两次运行产出不同报告。按原候选索引重排后修复——凡是并发产出的结果,都要问一句”顺序从哪来”。
- 心跳刷新失败被静默吞掉。吞掉它的任务会一直跑下去但租约早已过期,随时可能被别的 Pod 接管重跑,产生两个进程同时写同一个任务的竞态。心跳持久化失败必须让任务失败。
- 403 泄露任务存在性。越权返回 403 让任务 ID 枚举成为可能,统一 404 后才闭合。
- stale 阈值与心跳间隔倒挂。接管阈值必须显著大于心跳间隔,否则网络抖动一下任务就会被误判为 stale 双跑。两者是配套参数,不能各自独立调。
十、局限与适用边界
这套方案不是银弹,边界要自己先说清楚:
- 吞吐天花板是真实的。CAS 领取依赖行锁竞争和轮询,任务量级到每秒百条以上就该认真考虑专用队列了。当前场景离这个量级有两个数量级以上的余量。
- 轮询延迟是秒级。对”任务创建后多久开始执行”有亚秒级要求的场景,轮询模型不合适。
- 单任务内并发上限不等于全局限流。多 Pod 乘以任务内并发才是总在途请求,如果 LLM 供应商有严格的并发配额,这里只有软约束。
- 没有任务优先级、延时队列、死信队列这些 MQ 常见设施。当前不需要,需要时再评估是扩展任务表还是引入 MQ,而不是提前为想象中的需求建复杂度。
- 幂等是前提。租约接管意味着任务可能被重跑,流水线必须写成重跑安全的。如果执行的是不可重复的外部副作用,需要额外的去重层。
十一、总结
这一篇把视线从”单个任务内部怎么跑”拉到了”任务之间怎么被编排”。核心结论是:在中等吞吐下,PostgreSQL 事务加 CAS 领取加租约恢复,用一张任务表就撑起了长时异步任务的全部可靠性需求——原子分配、故障接管、优雅关闭、状态可查——同时把部署拓扑压到最简:一个数据库、一个应用镜像、一条本地启动命令。消息队列是伟大的工具,但它解决的问题(高吞吐、多消费组、削峰、跨语言)在诊断场景里一个都不存在,为不存在的需求引入真实的运维成本,是架构上典型的负收益。
回扣系列主线:模型负责语义判断,应用负责协议边界。任务编排就是”协议边界”在执行层的体现——谁拥有任务、谁可以读它、它挂了由谁接管、并发完成顺序如何不污染结果,这些都不是模型能兜底的,只能靠应用层的机制硬性保证。前面四篇解决了”一次诊断怎么跑对”,这一篇解决了”很多次诊断怎么跑稳”。
下一篇是系列的最后一篇:诊断明细堆在数据库里只是中间产物,怎么把它变成 PM 真正会看、看得懂、能据此决策的报告——从字段取舍到信息层级的产品化过程。