AI编程 擎天说

把 PostgreSQL Outbox 安全送进 JetStream:一次围绕故障恢复的落地复盘

这篇写了什么

把 PostgreSQL Outbox 安全送进 JetStream:一次围绕故障恢复的落地复盘 当 NATS 在发布过程中断开时,最麻烦的情况不是明确失败,而是服务端可能已经收到消息,客户端却没有拿到确认。此时如果直接重试,就可能产生重复消息;如果不重试,Outbox 又可能长期积压。 Happy

把 PostgreSQL Outbox 安全送进 JetStream:一次围绕故障恢复的落地复盘

当 NATS 在发布过程中断开时,最麻烦的情况不是明确失败,而是服务端可能已经收到消息,客户端却没有拿到确认。此时如果直接重试,就可能产生重复消息;如果不重试,Outbox 又可能长期积压。

HappyNM IM 这部分工作的目标,是把已经写入 PostgreSQL 的事务 Outbox 发布到 NATS JetStream,并在确认未知、消息总线停机、进程重连这些场景下保持至少一次发布语义,同时让重复消息能够稳定收敛。推进过程中,我借助 AI 开发工具先梳理需求和边界,再把实现拆成组件、故障门和检查项,而不是把它当成一个“补一个发布函数”的任务。

先把问题从“发送消息”改成“处理不确定性”

最初需要确认的并不是 SDK 怎么调用,而是跨 PostgreSQL 和 JetStream 时,哪些事情可以保证,哪些不能保证。

事务 Outbox 只能说明事件已经持久化。它无法和 JetStream 的发布确认组成一个跨系统事务,所以这里不把目标写成恰好一次。更实际的方案是:允许确认未知后的重复发布,再依靠 JetStream 的服务端去重窗口把同一个事件收敛成一条消息。

因此,事务 Outbox 的 eventId 被用作稳定消息 ID。这个选择让重试有了明确的身份依据,也避免每次重发都生成一个新的消息标识。AI 工具在拆解需求时帮助我把“安全发布”具体化为几个判断:事件信封是否固定、消息 ID 是否稳定、重复发布由谁处理、最终投递状态是否和用户已读状态混在一起。

最终,发布器使用严格版本化的 NATS JetStream Publisher,并生成 canonical 的内部事件信封。消息总线或 WebSocket 发布成功,只能说明服务侧完成了相应动作,不能直接推导出 delivered 或 read。用户回执仍由独立的权威游标推进,断线客户端则继续通过现有 Sync/Timeline 补缺。

给 lease 加一条不依赖 SDK 的时间边界

Outbox Worker 采用 PostgreSQL lease relay,从数据库领取待发布事件,发布完成后再更新状态。这里有一个容易被 SDK 行为掩盖的风险:发布 Promise 在连接重连期间可能持续 pending。

如果只依赖 SDK 自带的 publish timeout,数据库 lease 就没有一个可以证明的释放上界。Worker 可能一直占着这条记录,后续实例无法及时接管,积压也就不再只是消息总线的问题,而会扩散到数据库调度层。

我把这个风险单独交给 AI 工具分析,并将它拆成“连接断开时 Promise 如何结束”“lease 何时能被其他 Worker 重新领取”“停机与超时谁先生效”几个问题。实现上,在 SDK 之外增加独立的硬 deadline,并把它和 lease 时限、发布超时、优雅停机组合约束。这样,即使底层客户端在重连期没有按预期返回,Worker 仍然有明确的退出路径。

Worker 还使用 CSPRNG 生成 lease ID,支持优雅停机,提供 readiness 和 Prometheus 指标。它的职责保持独立:负责从 Outbox 领取并发布事件,不把 stream 创建、数据库变更等基础设施操作藏进运行流程里。

真实故障门暴露了两个不能靠演示发现的问题

第一轮设计里,停机测试使用的是匿名宿主端口。容器重启后端口会变化,结果是故障注入同时改变了客户端配置地址,测试无法只验证“NATS 暂时不可用”这一件事。

这不是业务代码本身的错误,却会让测试结论失真。我把故障门改成预留并显式映射稳定本地端口。这样,容器重启只改变 NATS 的可用性,客户端仍然连接同一个地址,停机和恢复场景才具备可重复的边界。

第二个问题出现在真实重连门。测试发现,SDK 自带 timeout 不能单独约束重连期间的 pending publish。也就是说,表面上配置了超时,实际 lease 仍可能没有按预期恢复。

补上外层硬 deadline 后,Worker 能够在发布确认未知时释放或结束当前处理,后续再由 Outbox 重新领取。重试使用原始 eventId 作为稳定消息 ID,JetStream 的去重机制负责把重复发布收敛。这里验证的不是“永远没有重复调用”,而是系统在无法确认结果时,仍能恢复并把最终消息数量控制在预期范围内。

启动检查只读基础设施,避免运行进程偷偷改状态

另一个判断是,Worker 启动时不应该顺手创建或修改基础设施。它只检查 migration 状态,以及已有 stream 是否满足兼容性要求;不自动修改数据库,也不自动修改已有 stream。

这样做牺牲了一点启动时的“自动方便”,换来的是更清楚的运维边界。基础设施变更由明确流程负责,运行进程只消费已经准备好的资源。对于消息流名称、版本和配置不匹配的情况,启动检查应该尽早暴露问题,而不是让 Worker 带着隐式变更继续运行。

这部分也被纳入 AI 辅助拆解:哪些检查属于启动前置条件,哪些动作会改变外部状态,哪些错误必须让 readiness 失败。最后形成的是独立的 stream provision 与兼容性检查,而不是把 provision 藏在发布路径中。

用组合测试确认恢复路径,而不是只看单元测试

落地检查分成几层。完整工程门禁中,82 个测试文件、799 项测试通过,类型检查、Lint、格式检查和构建也全部通过。PostgreSQL 18.3 与 17.10 的 Adapter 集成测试各 52 项保持通过,官方生产依赖审计返回 0 个已知漏洞;新增生产依赖许可证为 Apache-2.0 与 Unlicense。

更关键的是两项真实 Outbox Worker 组合故障用例:确认未知后重复发布,JetStream 最终只保留一条消息;NATS 停机后,原 Worker 能够重连并自动补发。这两项测试把稳定消息 ID、lease 恢复、外层 deadline 和 JetStream 去重放在同一条路径上验证,而不是分别证明几个孤立函数能工作。

这次推进里,AI 开发工具最有价值的地方不是代替我写完发布器,而是帮助我把模糊的可靠性要求拆成可反驳的假设,再根据故障门结果返工。端口映射问题说明测试环境也需要被检查;SDK 重连问题则说明配置项的名字不能代替行为验证。

当前阶段完成的是 Outbox 到 JetStream 的可恢复发布基础。下一检查点还包括 durable Realtime consumer、受会话保护的 WebSocket fanout、ack/backpressure 和撤权。发送方回执投影保持独立检查点,断线客户端继续依赖现有 Sync/Timeline 补缺。