消息与事件的状态模型:发送成功为什么不是业务完成
从命令、事件、通知和任务的语义出发,拆开生产、Broker 接收、消费、业务提交与确认五个彼此独立的完成点。
RabbitMQ 的可靠性说明把 publisher confirm 与 consumer acknowledgement 明确视为两次独立的所有权转移;前者不知道消费者是否处理,后者也不知道生产者如何发布。参见 RabbitMQ Reliability Guide。
先区分消息载体与业务语义
消息是跨进程传输的数据单元,事件、命令、通知和任务则是它承载的业务意图。事件描述已经发生且不可撤销的事实,例如“订单已支付”;命令要求某个明确能力执行动作,例如“为订单预留库存”;通知只提示接收方重新读取权威状态;任务描述可调度、可重试的工作。四者如果共用一个模糊的 type=UPDATE,消费者就无法判断重复到达时应该覆盖、拒绝还是再次执行。
事件应使用过去式并携带事实发生时的业务版本。命令必须有目标边界、截止时间和幂等键;通知必须允许丢失后通过权威查询恢复;任务需要记录尝试次数、下一次执行时间和终止条件。技术 Topic 只是运输分组,不能替代这些契约。真正稳定的设计先写出“不变量”,再决定路由键和 Broker。
五个完成点不能压缩成一个布尔值
一次跨服务动作至少经过本地业务提交、消息持久化或发布、Broker 接收、消费者业务提交、消费位点推进五个完成点。HTTP 返回成功通常只证明本地服务接受了请求;发送回调成功只证明 Broker 在特定配置下接收了记录;消费回调返回只证明代码路径结束;ack 或 offset commit 只证明这条记录不再按原位置重投。业务结果是否成立,必须由目标服务的权威状态证明。
最危险的实现是先确认消息再提交数据库。进程在两者之间崩溃时,Broker 认为记录已经完成,数据库却没有结果。反过来,先提交数据库再确认,崩溃会导致重投,但只要消费端幂等,重复是可恢复的。因此可靠消费通常有意选择“可能重复但不静默丢失”,而不是追求一个无法跨资源兑现的恰好一次口号。
事件信封是演进边界,不是字段垃圾桶
信封至少要稳定表达 eventId、eventType、occurredAt、producer、schemaVersion、correlationId、causationId、partitionKey 与 payload。eventId 用于去重和审计;correlationId 串起一次业务旅程;causationId 指向直接原因,防止把整条链误认为一个同步调用。partitionKey 决定局部顺序和负载分布,不能发布后随意改变。
不要把用户姓名、完整 URL、异常堆栈和无限增长的 headers 塞入信封。敏感数据会扩散到日志、重试队列和归档;大 header 会降低批处理密度。信封只承载跨边界需要的最小事实,详情通过受权查询获得。
用所有权转移图读懂可靠性
可靠消息不是一个开关,而是一组有证据的所有权转移。生产者在本地事务提交前拥有业务意图;发布成功后 Broker 对已接收记录承担保存和投递责任;消费者取得记录后承担处理责任;目标数据库提交后才形成新的业务事实;最后的 ack 或 offset commit 释放 Broker 的重投责任。任一阶段超时,都要区分“动作未发生”“动作已发生”和“结果未知”。
确认信息只能回答所属边界的问题。发布确认不证明消费完成,消费确认不证明所有外部副作用成功,队列为空也不证明业务状态正确。故障分析要保存 eventId、业务 key、生产时间、Broker 位置、投递次数、消费者实例、业务事务 id 和确认时间,才能重建状态迁移。
失败矩阵比“保证不丢”更可执行
工程选择的核心是把不可避免的不确定窗口转成可重复、可查询、可补偿的状态。不要用无限重试掩盖未知结果,也不要用提前确认换取表面低 lag。
从订单支付推导一条完整证据链
支付服务提交订单版本 7 和 Outbox 事件 evt-7 后,调用方得到的成功只能解释为支付事实已在本地成立。Relay 发布 evt-7 并得到 Broker 接收证据;库存消费者以 evt-7 去重,在库存事务中写入预留结果和消费记录;确认丢失导致 evt-7 再次投递时,消费者返回已保存结果而不重复扣减。通知消费者可以失败并独立重试,因为它不改变支付事实。
若库存长期不足,系统产生“库存预留被拒绝”这个新事实,而不是篡改或删除“订单已支付”。协调者据此发起退款补偿。这样每个事件都是可审计的过去事实,每个命令都有目标和终止状态,故障恢复通过追加状态推进,而不是试图把分布式历史回滚成从未发生。
用两个 Java 状态模型固定不变量
下面的模型不连接真实 Broker。它们把本篇最容易混淆的状态、位置或所有权变化压缩成确定输出;随后再用 RabbitMQ、Kafka 或 RocketMQ 的集成环境验证协议、持久化、重平衡和故障时序。
javac --release 17 -Xlint:all -Werror examples/backend-development/message-event/message-event-model/MessageCompletionDemo.java examples/backend-development/message-event/message-event-model/CommandEventSemanticsDemo.java
java -cp examples/backend-development/message-event/message-event-model MessageCompletionDemo
java -cp examples/backend-development/message-event/message-event-model CommandEventSemanticsDemopublished=true brokerAccepted=true processed=true businessCommitted=false acknowledged=false
commandTarget=inventory eventFact=order-paid notificationHint=true taskRetryable=true模型输出变化时,应定位是哪条业务不变量被修改,而不是只更新示例字符串。真实集成测试还要注入进程崩溃、响应丢失、Broker 切换和下游超时,因为这些窗口无法由纯内存模型模拟。
契约、权限与数据生命周期不能交给默认值
Topic、exchange、queue 和 consumer group 名称是运行拓扑,eventType 与 schemaVersion 才是业务契约。一个 Topic 可以承载多个兼容事件,也可以按敏感级别、保留期和吞吐隔离成多个 Topic;拆分依据应是治理边界,而不是每新增一个 Java 类就创建资源。契约目录要记录 owner、生产者、消费者、分区键、顺序承诺、投递语义、最大尺寸、保留期、兼容策略与废弃流程。
生产者升级时先发布消费者能够忽略或默认的新字段,再逐步升级消费者,最后停止旧字段;破坏性变化使用新 eventType 或新版本并经历双写、双读、回填和观测窗口。消费者不能因为遇到未知可选字段就失败,也不能把未知枚举静默映射为某个现有业务状态。无法理解的语义应进入受控隔离并告警,而不是无限重试。
用分层测试证明协议与业务同时成立
纯函数或状态模型验证 eventId、版本、重试和状态迁移;容器化集成测试验证真实客户端的 confirm、ack、offset、事务和拓扑参数;故障测试在消息已写出、业务已提交、确认未到达等精确位置杀进程;容量测试使用接近真实的消息尺寸、key 分布、批次和下游延迟。四层测试不能互相替代。
测试断言应落在权威结果:同一个 eventId 重放十次仍只有一次业务写入,旧版本不能覆盖新版本,毒消息不会阻塞健康 key,消费者撤销所有权后不能继续提交,积压恢复不会突破数据库连接预算。只断言“发送方法未抛异常”或“最终收到一条消息”,无法覆盖最重要的失败窗口。
把设计决定写成可以被证伪的承诺
指标、日志与容量要围绕状态迁移
eventId、订单号和用户号不能成为时序指标标签,否则基数会随消息量增长。它们进入结构化日志和采样 trace;指标只按 topic/queue、consumer group、result、failure class、schema version 等受控维度聚合。日志中的 payload 默认脱敏,死信查看需要权限、审计和保留期限。
容量模型至少计算平均/峰值每秒消息数、平均/p99 字节、保留时间、副本倍数、重试放大、批处理密度、消费者处理时长与下游并发。恢复测试要在稳定生产流量仍存在时注入积压,验证净排空速率,而不是暂停生产后测一个理想峰值。
上线与演练以可恢复为验收标准
上线前固定契约兼容测试、生产确认、消费幂等、毒消息隔离、死信修复、优雅停机和积压恢复。灰度按独立 consumer group 或路由键切流,避免新旧消费者竞争同一业务状态却使用不同语义。回滚不仅回滚代码,还要处理已经发布的新版本事件和新建拓扑。
故障演练依次注入生产响应丢失、Broker 节点切换、消费者提交后崩溃、下游超时、永久坏消息和热点 key。每次演练都回答三件事:业务不变量是否保持、重复或延迟是否有持久证据、系统是否在容量预算内自动或人工收敛。
