rocketmq事务消息源码解析

发布时间:2026/8/25 11:12:46
rocketmq事务消息源码解析 什么是事务消息事务消息是apache rocketMq提供的一种高级消息类型支持在分布式场景下保障消息生产和本地事务的最终一致性代码示例事务消息处理流程生产者将消息发送至rockerMq服务端rockerMq服务端口将消息持久化后向生产者返回ACK确认消息发送已经成功。此时消息被标记为“暂不能投递”这种状态下的消息即为半事务消息生产者开始执行本地事务逻辑生产者根据本地事务执行结果向服务端提交二次确认结果服务端接收到确认结果后处理逻辑如下二次确认结果为commit服务端将半事务消息标记为可投递并投递给消费二次确认结果为Rollback服务端将回滚事务不会将半事务消息投递给消费者若服务端未收到发送者提交的二次确认结果或服务端收到的二次确认结果为Unknow未知状态经过固定时间后服务端将对消息生产者即生产者集群中任一生产者实例发起消息回查生产者收到消息回查后需要检查对应消息的本地事务执行的最终结果生产者根据检查到本地事务的最终状态再次提交二次确认。服务端仍按步骤4对半事务消息进行处理底层逻辑处理流程基础流程事务回查数据流转图关键源码1、事务消息发送2、half消息存储broker存储half消息3、执行本地事务提交事务状态4、处理事务状态5、事务回查事务回查时间间隔事务状态提交超时时间底层逻辑触发回查的条件条件含义opMsg null valueOfCurrentMinusBorn checkImmunityTime没有对应的 Op 消息且已超过免疫期immunity time确保消息有足够时间等待生产者提交状态不至于太早回查。超时就需要回查opMsg ! null opMsg.getLastBornTimestamp() - startTime transactionTimeoutop 消息存在但时间戳太早可能未覆盖当前消息。防止op记录延迟因未记录导致事务消息不处理valueOfCurrentMinusBorn -1消息时间异常未来时间举例源码定时任务进行事务回查生产者处理回查请求存储设计half和op两个队列组成部分作用设计原因half队列RMQ_SYS_TRANS_HALF_TOPIC存储事务的“准备”消息半消息未提交前不可见保证消息先持久化防止事务状态丢失op 队列RMQ_SYS_TRANS_OP_HALF_TOPIC存储事务的“操作”消息commit/rollback记录事务状态变化实现幂等性与恢复机制OP消息Op 消息体是一个字符串内容格式如下offset1[offsetSeparator]offset2[offsetSeparator]deleteContext缓存每个 queueId 对应一个 MessageQueueOpContext内部维护一个队列和总数据大小批量写入op队列批量写入时间间隔写入逻辑遍历deleteContext构建OP消息写入commitLog根据queueId取出deleteContext中的队列contextQueue,进行追加为消息体构建OP消息事务消息的状态恢复机制当 Broker 故障重启后通过 Op 消息可以重建事务状态Op 消息记录了哪些 half 消息已经被提交或回滚在 Broker 启动时会加载这些 Op 消息重建 removeMap清理已处理的 half 消息对于未处理的消息后续会触发回查check机制。事务消息的幂等性唯一业务 IDUNIQ_KEY通常为 msgId 或业务自定义 IDOp 消息中的 queueOffset 列表确保每个事务只处理一次生产者本地事务幂等性逻辑避免重复执行本地事务。