RocketMQ 运行链:CommitLog、消费队列与事务回查的边界
拆解 CommitLog、ConsumeQueue、普通与顺序消息、重试以及事务消息回查,避免把半消息误解为分布式事务。
RocketMQ 官方 Transaction Message 描述半消息、本地事务与事务状态回查;FIFO Message 把顺序限定在同一消息组。
CommitLog 与 ConsumeQueue 分离了顺序写和按主题消费
Broker 先把消息顺序追加到 CommitLog,再为 Topic 与队列构建较轻的 ConsumeQueue 索引。索引保存物理偏移、大小和标签摘要,消费者先读索引再定位消息体。这个结构解释了为什么顺序写可以获得高吞吐,也解释了磁盘抖动、索引构建滞后和大消息会如何影响读取链。
队列是并行消费和局部顺序的基本单位。相同业务实体需要稳定映射到同一消息组或队列;如果 key 分布倾斜,一个热点订单族会把整个队列拖慢。增加队列数不会自动重排历史业务顺序,发送选择器与消费者并行度必须一起演进。
事务消息协调可见性,不接管本地事务
事务生产者先发送对消费者不可见的半消息,Broker 接收后回调生产者执行本地事务。生产者随后提交或回滚消息;若 Broker 未得到确定结果,会向生产者回查本地事务状态。可靠性的关键不是回调返回值,而是生产者能从本地权威记录中确定回答 COMMIT 或 ROLLBACK。
因此本地事务必须把业务状态与可回查标识一起提交。若回查代码只看进程内变量,重启后就会永远 UNKNOWN。回查还应幂等、快速且可限流,避免 Broker 扫描未决事务时反压业务数据库。事务消息只解决“本地事务成功后消息最终可见”,下游消费仍可能重复、失败和进入重试。
顺序、重试与可用性需要显式交换
FIFO 消息按消息组提供顺序。前一条持续失败时,后续同组消息不能越过,否则状态机顺序被破坏;这意味着毒消息会阻塞整个组。业务必须决定是暂停并修复、把失败转为补偿状态,还是允许在明确规则下跳过。顺序不是免费属性,而是把并行度和故障隔离缩小到一个组。
发送重试处理客户端到 Broker 的暂态失败,消费重试处理业务执行失败,两者具有不同的幂等风险。超时后结果不明时,生产端重试可能产生重复;消费失败时,Broker 重投也可能发生在业务已经提交之后。统一 eventId 与本地状态约束是两端共同的恢复锚点。
用所有权转移图读懂可靠性
失败矩阵比“保证不丢”更可执行
工程选择的核心是把不可避免的不确定窗口转成可重复、可查询、可补偿的状态。不要用无限重试掩盖未知结果,也不要用提前确认换取表面低 lag。
消费模式改变所有权与可见性窗口
拉取模式由客户端维护拉取位置和消费进度,适合批量、流控和精细位点管理;Pop 模式更接近不可见时间窗口,消息投递后暂时对其他消费者隐藏,处理完成后确认,超时未确认则再次可见。不可见时间必须大于正常处理尾延迟,又不能大到故障后长时间无法重投。
业务处理可能超过单次不可见窗口时,应续期或拆分任务,但续期也需要截止时间,不能让僵死消费者永久占有消息。无论哪种模式,客户端看到成功确认都只释放 Broker 责任;事务回查、消费重试和死信分别解决不同阶段,不能用一个事务消息类型覆盖整条链。
用两个 Java 状态模型固定不变量
下面的模型不连接真实 Broker。它们把本篇最容易混淆的状态、位置或所有权变化压缩成确定输出;随后再用 RabbitMQ、Kafka 或 RocketMQ 的集成环境验证协议、持久化、重平衡和故障时序。
javac --release 17 -Xlint:all -Werror examples/backend-development/message-event/rocketmq-runtime/CommitLogConsumeQueueDemo.java examples/backend-development/message-event/rocketmq-runtime/TransactionCheckDemo.java
java -cp examples/backend-development/message-event/rocketmq-runtime CommitLogConsumeQueueDemo
java -cp examples/backend-development/message-event/rocketmq-runtime TransactionCheckDemocommitLogOffset=128 consumeQueueEntry=orders:2->128 bodyLoaded=true
halfMessage=true localState=COMMITTED checkResult=COMMIT visible=true downstreamRetryable=true模型输出变化时,应定位是哪条业务不变量被修改,而不是只更新示例字符串。真实集成测试还要注入进程崩溃、响应丢失、Broker 切换和下游超时,因为这些窗口无法由纯内存模型模拟。
契约、权限与数据生命周期不能交给默认值
用分层测试证明协议与业务同时成立
把设计决定写成可以被证伪的承诺
选型也应落到这些约束。RabbitMQ 的灵活路由与队列语义适合任务分发和复杂路由,Kafka 的分区日志与保留适合事件流、回放和多消费组,RocketMQ 的消息类型与事务回查适合其生态中的业务消息。产品能力不同,但任何一个都不会替应用定义业务事实、提交外部数据库或处理不可逆副作用。先完成语义和失败矩阵,再选择实现,系统才不会被某个客户端注解反向塑形。
指标、日志与容量要围绕状态迁移
容量模型至少计算平均/峰值每秒消息数、平均/p99 字节、保留时间、副本倍数、重试放大、批处理密度、消费者处理时长与下游并发。恢复测试要在稳定生产流量仍存在时注入积压,验证净排空速率,而不是暂停生产后测一个理想峰值。
上线与演练以可恢复为验收标准
故障演练依次注入生产响应丢失、Broker 节点切换、消费者提交后崩溃、下游超时、永久坏消息和热点 key。每次演练都回答三件事:业务不变量是否保持、重复或延迟是否有持久证据、系统是否在容量预算内自动或人工收敛。
