RabbitMQ 可靠消息实战:从同步 HTTP 到异步 MQ,Outbox/Inbox、幂等与灰度迁移设计
需求编号:DMS-1285
分支:feature/DMS-1285
实现提交:df7f43471
整理日期:2026-08-09
前言
表面上看,这个需求只是把 DMS 随附单据的敏感词校验从同步 HTTP 调用改成 RabbitMQ 异步调用。但真正开始设计后会发现,协议替换反而是最简单的一部分。
真正困难的是:生产环境中已经存在大量按旧逻辑处理的单证,MQ 又必须分批开放;一个单证可以上传多个附件,也可以删除、重复上传和重新上传;消息可能重复、延迟或丢失;审核结果回来时,原附件甚至可能已经不存在;最重要的是,HTTP 和 MQ 两条链路必须长期共存,但同一个业务对象绝对不能在两条链路之间来回切换。
因此,这次实现解决的并不是单纯的“异步化”,而是一组彼此关联的问题:如何确定并锁定传输链路,如何可靠地发布和消费消息,如何描述多附件的聚合状态,如何让历史确认在附件变化后自动失效,以及如何在不要求全国业务人员重新上传附件的前提下逐步扩大全量客户。
为什么这个 RabbitMQ 改造案例值得写成一篇博客
并不是每一个业务需求都值得单独写一篇技术文章。很多需求只是在既有模式中增加字段、页面或判断条件,它们对具体项目很重要,但脱离业务背景后未必还有多少可复用价值。DMS-1285 不同,它集中体现了老系统向异步架构演进时最常见、也最容易被低估的一组工程问题。
首先,它不是一个可以推倒重来的“新系统案例”,而是在已有生产数据、既有 HTTP 链路和全国业务人员持续操作的前提下完成迁移。真实系统的难点通常不在于设计一条理想的新链路,而在于如何让新旧规则共存,既不破坏历史行为,也不给用户增加补数据和重新上传的负担。这种渐进式改造经验,比单独介绍 RabbitMQ 的使用方法更接近企业系统的实际处境。
其次,这个需求包含了多个具有普遍性的架构矛盾:可用性与一致性如何取舍,灰度开关与业务路由如何分离,数据库事务与消息发布如何保持最终一致,多附件状态如何聚合,重复和迟到消息如何幂等处理,以及人工确认如何绑定数据版本。它们并非 DMS 独有,在订单、支付、风控、审批、文件处理等系统中都会反复出现。
再次,方案形成过程本身具有记录价值。最初直观的创建时间切流、MQ 失败后回退 HTTP、用任务是否存在代替路由等方案,看起来都能少写一些代码,但继续推演历史数据、二次扩面和异常恢复后,都会暴露出新的风险。把这些被否决的方案和理由写出来,比只展示最终代码更能说明设计是如何收敛的,也能帮助后来者少走弯路。
最后,这个需求形成了一个完整的工程闭环:前端提示、服务端守卫、不可变路由、文件状态机、Outbox/Inbox、结果消费、超时处理、死信告警、灰度开关和上线策略彼此配合。它说明“接入 MQ”不只是成功发送一条消息,而是要让整个业务在正常、失败、重复、延迟、删除和扩容等状态下都有可解释的行为。
所以,这篇博客真正值得记录的不是改动了几十个文件或几千行改动,而是一套可以迁移到其他存量系统的思考方法:先确定不能被破坏的业务不变量,再选择技术方案;先设计新旧链路如何共存,再讨论如何切换;先把失败和重试纳入正常流程,再谈异步化带来的性能收益。
一、为什么要从同步 HTTP 迁移到 RabbitMQ
一期逻辑主要依赖同步 HTTP:在业务动作发生时,DMS 直接请求审核平台并等待结果。这种方式结构直观,但存在几个天然限制。
首先,审核平台的可用性和响应时间会直接影响 DMS 请求线程。网络抖动、平台超时或瞬时不可用,都可能放大为业务页面卡顿或操作失败。
其次,同步调用很难完整表达附件生命周期。一个单证可能有多个文件,文件可能被删除或重新上传,后续结果必须精确对应到“当时的那一个文件”,而不能只依靠文件名或单证级状态。
最后,升级不能一刀切。生产环境中的历史单证仍应保持一期行为,新创建或首次进入新链路的业务才适合逐步走 MQ。如果只设置一个全局切换时间,很容易在后续扩大客户范围时遇到补建任务、重新上传附件或链路混用的问题。
二、HTTP 迁移 RabbitMQ 时会遇到哪些核心问题
1. HTTP 和 MQ 如何共存:为什么不能失败后 fallback
系统必须同时保留两种传输方式:历史业务继续走 HTTP,新业务按开关进入 MQ。然而“失败后从 MQ fallback 到 HTTP”在这里并不是容灾,而是破坏一致性。
如果同一个单证先向 MQ 发过请求,超时后又调用 HTTP,审核平台可能收到两份请求;两份结果的先后顺序不可控,DMS 也无法判断哪一份才是最终结果。反过来,如果历史 HTTP 单证因一次重新上传突然进入 MQ,也会改变生产中已经形成的业务语义。
因此最终规则是:链路一旦确定便永久锁定。MQ 失败仍在 MQ 链路内重试和告警,HTTP 链路也只使用 HTTP,绝不跨链路兜底。
2. RabbitMQ 灰度切流:为什么不能只按“单证创建时间”判断
创建时间切流适合一次性全量上线,但不适合“先验证 FNSR,再逐步扩大全客户”的场景。
假设首次上线时只开放 FNSR,一个非 FNSR 单证在此期间已经创建并上传附件。以后打开全量开关,如果仍按创建时间判断,它可能永远留在旧链路;如果修改切换时间,又可能让已经处理过的单证突然改走 MQ。为了补齐状态,只能要求用户重新上传或由后台批量补任务,操作成本和风险都很高。
因此,时间不是可靠的业务边界。真正稳定的边界应当是“这个单证第一次需要选择校验链路时,现场已经存在什么事实”。
3. 多附件场景:为什么必须使用文件级状态而不是单证级状态
敏感词审核的最小实体必须是文件,而不是单证。单证只负责聚合多个文件的状态。
如果只保存单证级结果,就无法正确处理以下情况:两个同名文件、删除其中一个文件、删除后重新上传、旧文件结果迟到、一个文件正常而另一个文件待人工审核。文件名也不能作为唯一标识,因为重复文件名在业务上是允许出现的。
4. 数据库事务与 RabbitMQ 如何避免双写不一致
如果先提交业务数据再直接发 MQ,进程可能在两步之间崩溃,形成“数据库有任务、MQ 没消息”;如果在数据库事务中先发 MQ,事务随后回滚,则会形成“MQ 已收到消息、数据库却没有任务”。这就是典型的双写不一致。
5. RabbitMQ 消息重复、延迟和乱序如何处理
生产者重试、网络恢复和消费者重连都可能造成重复投递。结果也可能在附件删除后才返回,甚至在任务超时后才到达。消费者若把“收到消息”直接等同于“业务处理成功”,就会产生重复更新或丢结果。
6. 附件变化后为什么历史人工确认必须失效
用户确认过一次风险结果,并不代表以后上传的新版本附件也已确认。附件删除、重新上传、重新校验或审核结果变化后,旧确认必须失效,否则会留下绕过审核的漏洞。
三、RabbitMQ 可靠消息整体架构:路由、任务与消息可靠性分层
最终实现把问题拆成四个持久化模型:
ATTACH_DOC_CHECK_ROUTE:记录一个业务维度永久选择LEGACY还是MQ。ATTACH_DOC_CHECK_FILE:记录每个附件的校验任务、请求标识、状态和结果。ATTACH_DOC_MQ_OUTBOX:记录待发布、发送中、已发送或待重试的请求消息。ATTACH_DOC_MQ_INBOX:记录已经消费的结果消息,实现幂等。
主链路可以概括为:
附件上传
-> 判断是否属于敏感词校验范围
-> 查询或创建不可变路由(LEGACY / MQ)
-> LEGACY:保持原有 HTTP 行为
-> MQ:同一事务保存文件任务和 Outbox
-> 后台发布器投递请求消息
-> 审核平台消费并返回结果
-> DMS 结果消费者通过 Inbox 幂等落库
-> 页面、审核和发送入口读取文件聚合状态
这里最重要的设计原则是:路由只回答“走哪条链路”,文件任务只回答“这个文件现在是什么状态”,Outbox/Inbox 只回答“消息是否可靠发布和可靠消费”。三个问题不再相互混用。
四、HTTP 与 RabbitMQ 灰度迁移:用不可变路由兼容历史业务
4.1 HTTP / MQ 路由应该按什么粒度锁定
路由表以“报关单 + 单证类型 + 关联标识”为唯一业务维度。无关联对象时统一使用 RELATION_KEY=0。数据库唯一约束避免并发情况下为同一个业务维度创建两条路由,触发器则禁止删除路由或修改已经确定的通道。
采用独立路由表,而不是给原来的三张业务表增加字段,有两个原因:一是多种单证来源可以统一表达,不必在多个旧表中重复迁移;二是路由是本需求新增的跨表业务事实,独立建模更容易施加唯一约束和不可变规则。
4.2 附件上传时如何选择 HTTP 或 MQ 路由
SensitiveWordRouteService.resolveForUpload 按以下顺序判定:
| 优先级 | 已知事实 | 路由 | 分配原因 |
|---|---|---|---|
| 1 | 已有路由记录 | 复用原路由 | 保证不可变 |
| 2 | 当前有效附件已有 MQ 任务 | MQ | EXISTING_MQ_TASK / 已存在有效 MQ 校验任务 |
| 3 | 存在历史附件 | LEGACY | EXISTING_ATTACHMENT / 已存在历史附件 |
| 4 | 无历史附件,MQ 已开启 | MQ | FIRST_UPLOAD / 首次上传且 MQ 已启用 |
| 5 | 无历史附件,MQ 未开启 | LEGACY | MQ_DISABLED / 首次上传时 MQ 未启用 |
上传完成后再次解析时,会排除“本次刚保存的文件记录”。否则第一份新附件刚入库就会被误认为历史附件,错误锁定为 LEGACY。代码还同时校验文件来源类型,避免不同附件表中数值相同的主键被误删或误判。
4.3 审核时如何复用 HTTP / MQ 路由
审核入口不主动制造新的 MQ 任务。它先复用已有路由;没有路由时,如存在有效 MQ 任务则锁定 MQ,如存在历史附件则锁定 LEGACY;完全没有附件时不创建路由。
这保证“上传”是任务产生点,“审核”只是状态校验点,不会因为用户打开一次审核页面就改变附件链路。
4.4 为什么不可变路由能支持 RabbitMQ 平滑扩面
默认情况下,仅 FNSR 订单进入敏感词校验。SWITCH_CONFIG 中的 sensitive_word.full_check 可以在 FNSR 验证稳定后打开全量校验。
打开开关后,不需要全国业务人员重新上传附件:已有附件的非 FNSR 单证在第一次进入相关操作时会因为“存在历史附件”锁定 LEGACY;尚无历史附件的新业务在首次上传时进入 MQ。HTTP 因此不需要在 MQ 上线后立刻关闭,它可以长期为 LEGACY 路由服务,直到历史业务自然消退。
这比时间切流更安全,因为它依据的是业务事实,而不是一个以后可能需要反复修改的时间参数。
五、Transactional Outbox:数据库与 RabbitMQ 如何保证最终一致性
MQ 路由下,SensitiveWordCheckService.createCheckTask 会为单个文件创建校验任务和 Outbox 记录,并在同一个 Cayenne 事务中提交。requestId 使用 UUID,同时写入 JSON 报文、AMQP message_id 和 correlation_id,成为贯穿请求、返回和幂等处理的统一标识。
实现过程中曾出现一个典型问题:outbox.setToAttachDocCheckFile(checkFile) 抛出空指针,而 checkFile 变量本身并不为空。调查后发现原因是 Cayenne 新对象尚未保存,数据库序列还没有分配主键,关系赋值时缺少可用对象标识。
最终顺序调整为:
save(checkFile)
save(outbox)
outbox.setToAttachDocCheckFile(checkFile)
commit
两个对象仍处于同一事务中,因此既获得了关系所需的对象标识,也没有牺牲原子性。
HTTP 上传线程在任务和 Outbox 提交成功后即可结束。后台 SensitiveWordRequestPublisher 每 10 秒扫描一次待发送数据,每批最多 100 条;查询直接在数据库中过滤 PENDING,以及已经到达下次重试时间的 RETRY,并使用 fetchLimit,避免先取大量数据再在 Java 内存中筛选。
发布过程使用持久化消息、mandatory return 和 publisher confirm。发布成功标记为 SENT,失败则保留同一个 requestId 进入 RETRY,不会创建一份新的业务任务。发布前还会确认源附件仍然有效;若附件已经删除,则取消对应的 Outbox 和文件任务,避免把无效文件发送给审核平台。
Outbox 解决的核心问题是:业务事务成功并不等于消息已经成功发布。只有把“待发送消息”本身也作为业务事务的一部分持久化,后台才有机会在进程重启或网络恢复后继续完成发布。
六、RabbitMQ 重复消费如何解决:Inbox + 手动 ACK 实现幂等
审核平台返回消息后,入口从 SensitiveWordResultConsumer 开始。消费者使用手动 ACK,并在独立 Cayenne DataContext 中处理结果:
- 校验 AMQP 标识和 JSON 报文格式;
- 以
messageId写入或查询 Inbox; - 按
requestId定位文件任务; - 忽略已经删除、失效或已处理的旧任务;
- 更新文件状态、原始结果、敏感词数量和低置信度数量;
- 数据库提交成功后才 ACK。
协议不合法的消息会被 reject 到结果死信队列;数据库异常则 NACK 并重新入队,同时中断当前 channel,避免异常消息在紧密循环中持续占用 CPU。
Inbox 的意义不是阻止 RabbitMQ 重复投递,而是让重复投递变得安全。网络层可以“至少一次”投递,业务层通过幂等记录得到“同一结果只生效一次”的效果。
七、多附件场景如何设计文件生命周期、聚合状态和版本化确认
7.1 文件级状态:删除、重传和迟到结果如何处理
每次上传以实际文件记录主键和来源类型建立任务。重新上传会产生新的文件身份;旧任务会失效,迟到结果不会覆盖新文件。删除附件时,相关任务和未发送 Outbox 同步失效。
这也是为什么不能用文件名做主键:多个同名文件、相同内容的重复上传和删除后重传都必须能被独立识别。
7.2 多附件状态如何聚合到单证级状态
页面、审核和发送校验不直接读取某一条文件记录,而是聚合当前有效附件。聚合优先级为:
PROCESSING > FAILED > REJECTED > MANUAL_REVIEW > NORMAL
当前附件没有对应任务时为 MISSING。其中 MISSING、PROCESSING 和 FAILED 会阻止审核或发送;正常状态放行;敏感词或低置信度结果要求人工确认。
校验不仅放在页面按钮上,也进入服务端审核和发送入口。这样即使绕过前端、通过其他页面或重试任务触发发送,也不能跳过敏感词校验。
7.3 人工确认为什么必须绑定当前附件版本
系统使用所有当前有效任务 requestId 排序后的 SHA-256 摘要作为 versionToken。人工确认节点保存:
SENSITIVE_WORD_CHECK:{versionToken}
再次校验时,不仅要求 token 相同,还要求确认时间不早于最新任务更新时间。只要附件被删除、重新上传、重新校验或任务版本发生变化,token 就会改变,历史确认自动失效。
这解决了一个容易被忽略的安全问题:用户确认的是“当前这一组审核结果”,而不是永久确认整个单证。
八、RabbitMQ 超时、重试、死信队列和后台任务如何设计
指定的 MQ worker 节点会启动一个包含两个线程的轻量调度池,承载三类周期任务:
- 每 10 秒派发 Outbox;
- 每 30 秒保证结果消费者和请求 DLQ 告警消费者处于运行状态;
- 每分钟扫描长时间未返回的 PROCESSING 任务。
使用两个线程是有意的:发布器等待 publisher confirm 时可能阻塞,如果所有任务共用一个线程,结果消费者维护和超时扫描都会被拖延。项目中的 Quartz、缓存刷新或其他业务线程池生命周期和阻塞特征不同,直接复用反而会扩大相互影响。
超时扫描只负责把超过配置时长、仍未返回的活动任务标记为 FAILED / RESULT_TIMEOUT,它不等同于把 MQ 消息送入死信队列。迟到结果如果仍对应有效文件且没有最终结果,仍可以恢复任务状态。
扫描性能的关键不在“每 10 秒”这个数字本身,而在查询是否走索引、是否限制批次。Outbox 已按状态、下次重试时间和创建时间建立索引,并限制每批 100 条;上线后仍应通过真实执行计划和表增长情况确认索引生效。
死信告警直接监听现有的 DMS:SENSITIVE_WORD_CHECK:REQUEST:DLQ,没有新建队列或绑定。消费到请求死信后调用统一的 DingTalkAlert,使用专属配置 dingtalk.mqdlx.webhook;DMS-1263 原有的 dingtalk.webhook 仍可通过同一工具类继续使用。只有钉钉 HTTP 成功且返回 errcode=0 才 ACK,发送失败则重新入队;相同内容默认有 5 分钟冷却时间,避免告警风暴。
需要明确:直接消费 DLQ 意味着告警成功后消息会从 Ready 列表移除。这符合“监听并处理”的语义,但不等于把死信永久保留在队列中。
九、RabbitMQ 可靠消息设计中为什么没有采用这些常见方案
| 备选方案 | 看起来的优点 | 实际缺点 |
|---|---|---|
| 按单证创建时间切流 | 配置简单 | 后续扩大全客户需要改时间、补任务或重新上传,容易误分历史单证 |
| MQ 失败后 fallback HTTP | 表面可用性高 | 同一业务产生双请求、双结果,链路语义被破坏 |
用是否存在 requestId 判断链路 |
少建一张表 | 任务缺失或清理后可能错误切换;任务身份不能代替路由事实 |
| 直接在事务后发布 MQ | 代码少 | 数据库提交与 MQ 发布之间存在崩溃窗口 |
| 在数据库事务内先发 MQ | 看似同步 | MQ 成功后数据库仍可能回滚,外部已看到不存在的业务 |
| 只保存单证级结果 | 表结构简单 | 无法正确处理多文件、删除、重传和迟到结果 |
| 以文件名去重 | 实现直观 | 同名文件和重复上传会碰撞,不能表达真实文件生命周期 |
| 只在前端按钮校验 | 改动小 | 其他发送入口、接口或后台路径可以绕过 |
| 为钉钉告警再建队列 | 隔离清晰 | 当前需求只需处理既有请求 DLQ,新增拓扑增加部署和运维成本 |
| 使用数字型路由原因码 | 存储短 | 排查时不可读;英文 code + 中文说明更适合日志、数据库和页面诊断 |
十、这套 RabbitMQ 可靠消息架构带来了什么收益
1. 老业务和新业务可以长期共存
历史附件自动保留 HTTP 行为,新业务进入 MQ;两条链路不依赖一次性停机切换,也不要求旧数据整体迁移。
2. 从 FNSR 扩大全客户不需要用户配合
全量开关只改变“尚未选路由的业务”是否纳入校验。已有附件按历史事实锁定 LEGACY,新业务自然进入 MQ,避免全国业务人员重新上传。
3. 消息失败可恢复、可追踪
Outbox 保存发送状态和重试时间,Inbox 保存消费幂等记录,文件任务保存业务结果。排查时可以区分“未发布、发布失败、等待结果、结果失败、业务已完成”,不再只依赖 RabbitMQ 管理台的瞬时曲线。
4. 多附件和重传场景有明确语义
每个文件独立校验,单证统一聚合;删除、重复文件、重新上传和迟到结果都有稳定的归属规则。
5. 人工确认不会越过版本边界
确认与当前任务集合绑定,文件变化后旧确认自动失效,减少误放行风险。
6. 运维闭环更加完整
请求死信可以主动推送钉钉,超时任务能在数据库中显式失败,发布端和消费端都有持久化状态,问题不再只能依靠人工盯队列。
十一、从 HTTP 灰度迁移到 RabbitMQ 的上线策略
建议按以下顺序上线,而不是设置“上线后一小时”的创建时间阈值:
- 先执行并核对四张业务表、索引、约束、触发器及
SWITCH_CONFIG脚本。 - 只在一个指定节点开启
dms.sensitive.mq.worker.enabled=true,其他节点保持关闭。 - 保持全量开关关闭,先验证默认 FNSR 范围;HTTP 继续为 LEGACY 路由服务。
- 观察 Outbox 积压、重试、结果延迟、DLQ、钉钉告警和数据库执行计划。
- FNSR 稳定后打开
sensitive_word.full_check,让尚未分配路由的非 FNSR 业务按附件历史自动选路。 - 不设置 fallback,不因单次 MQ 故障修改已确定的路由。
十二、RabbitMQ 可靠消息与灰度迁移的测试重点
代码级验证已覆盖 Java 8 Maven 编译,最近一次执行为 BUILD SUCCESS;git diff --check 也未发现空白错误。项目当前没有为本需求新增可执行的自动化测试,因此不能把“编译通过”表述为“全部场景已经自动化验证通过”。MQ 主链路和模拟审核平台已做过人工联调,但上线前仍建议按下表进行完整回归。
| 场景 | 重点断言 |
|---|---|
| 历史单证已有一个附件 | 首次进入时锁定 LEGACY,之后始终走 HTTP |
| 新单证首次上传一个附件 | MQ 开启时锁定 MQ,产生一个文件任务和一个 Outbox |
| 同一窗口上传多个附件 | 每个有效文件有独立任务,单证聚合状态正确 |
| 上传多个同名文件 | 不按文件名覆盖或错误去重 |
| 删除其中一个文件 | 对应任务失效,未发送消息取消,其他文件不受影响 |
| 删除后重新上传 | 产生新文件身份和新版本,旧结果不能覆盖新任务 |
| 重复投递同一结果 | Inbox 幂等,只处理一次 |
| 结果在删除后迟到 | 忽略无效旧任务,不影响当前聚合状态 |
| 任务先超时、结果后到 | 有效任务可以按规则恢复,不误关联其他文件 |
| MQ 发布失败 | 只进入 Outbox 重试,不调用 HTTP |
| MQ 请求进入 REQUEST DLQ | 钉钉成功后 ACK;钉钉失败时重新入队 |
| FNSR 开关阶段 | FNSR 进入校验,非 FNSR 不创建任务、不被拦截 |
| 打开全量开关 | 历史非 FNSR 附件锁定 LEGACY,新业务自然进入 MQ |
| 进口/出口、DEC/INV、关联单证 | 各入口使用相同路由和聚合规则,不遗漏来源类型 |
| 风险结果人工确认 | 当前版本可放行,附件变化后旧确认失效 |
| 绕过页面直接审核或发送 | 服务端守卫仍能拦截 MISSING、PROCESSING、FAILED |
十三、当前 RabbitMQ 方案的已知风险与后续工作
这套设计解决了链路一致性,但仍有几项生产风险需要正视:
- 当前调度模型要求只有一个节点开启 worker。多节点同时开启时,现有实现不能声称具备严格的跨 JVM 原子抢占,需要通过部署约束保证单 worker。
- Outbox、Inbox 和文件任务会持续增长,应根据审计要求设计归档或清理周期,并确保清理不会破坏仍有效的路由和幂等窗口。
- 发布器按单条 confirm 工作,每批最多 100 条。正常附件量下足够简单可靠,但高峰吞吐需根据真实监控评估。
- REQUEST DLQ 被告警消费者直接消费,成功告警后死信不再保留在 Ready 列表,运维侧需要接受并记录这一行为。
总结:老系统如何平滑迁移到 RabbitMQ
DMS-1285 最值得复用的经验,是不要把“升级传输协议”误认为“替换一个调用方法”。当系统已经有生产历史、多个附件生命周期、异步消息和分批放量要求时,最重要的是先定义不可变的业务边界,再分别解决消息发布、消息消费和业务状态问题。
路由表保证 HTTP 与 MQ 长期共存而不串线;Outbox/Inbox 解决跨系统消息的可靠性;文件级任务和版本化确认解决多附件生命周期;动态客户开关则让试点到全量的过程不需要用户补数据。每一个设计都比直接调用多了一点结构,但这些结构换来的,是上线时可解释、故障时可恢复、扩面时不扰民,以及结果真正可信。