
【RocketMQ】事务消息详解【一】什么是 RocketMQ 事务消息【二】使用场景【1】订单创建后扣减库存【2】支付完成发送通知、积分、优惠券【3】业务流程异步解耦需要强一致性【4】分布式链路中本地事务触发下游异步调用【三】底层原理【1】完整流程1发送半消息Prepare 消息2Producer 执行本地事务3Broker 事务回查核心容错机制【2】存储层面【3】时序图文字版【四】优缺点【五】替换方案分布式事务方案对比替代方案示例本地消息表【六】代码详细案例SpringBoot RocketMQ 事务消息【1】maven依赖【2】application.yml【3】实现 RocketMQLocalTransactionListener【4】生产者发送事务消息【5】OrderService 业务层【6】消费者库存服务注意消费端必须做幂等【七】注意事项【一】什么是 RocketMQ 事务消息RocketMQ 事务消息用于解决分布式事务问题实现本地数据库事务与消息发送的原子性要么本地事务执行成功 消息投递成功要么本地事务失败 消息完全不投递避免出现数据库成功消息没发 / 消息发出去数据库失败这种数据不一致。总结RocketMQ 事务消息保证生产端的事务和消息投递状态同步要么都成功要么都失败。注意RocketMQ 事务消息不能保证消费端事务原子性只保证「发送消息 本地事务」的原子性消费端还需要自己做幂等。【二】使用场景核心场景跨系统 / 跨服务本地 DB 操作要和发消息绑定原子性【1】订单创建后扣减库存订单服务创建订单本地事务同时发送消息给库存服务扣库存。如果订单入库成功消息没发出去 → 库存不减订单处于已创建无扣减数据不一致如果消息先发订单入库失败 → 库存扣了订单没生成。用事务消息订单创建成功消息对外可见订单回滚消息直接丢弃。【2】支付完成发送通知、积分、优惠券支付服务本地更新支付记录同时发消息给用户加积分、发优惠券、推送短信。支付成功才允许消息被消费支付失败消息不能被消费。【3】业务流程异步解耦需要强一致性业务 A 完成本地数据库操作后触发下游多个异步业务不能出现下游收到消息但上游数据库回滚。【4】分布式链路中本地事务触发下游异步调用比如用户注册成功初始化账户、生成账单多个下游服务异步处理。❌ 不适合场景1需要多服务数据库同时原子提交2PC 那种强一致性事务消息最终一致性不是强一致2短同步链路不需要异步3追求极低延迟事务消息有半消息、回查逻辑会有额外耗时。【三】底层原理RocketMQ 事务消息分为 半消息 (Prepare 消息) → 执行本地事务 → 提交 / 回滚消息 → 事务回查【1】完整流程1发送半消息Prepare 消息Producer 发送事务半消息给 BrokerBroker 收到消息标记为 TRANSACTION_PREPARED这条消息对消费者不可见消费者消费不到。返回成功给生产者。半消息消息已经存 broker但是不会投递消费队列。2Producer 执行本地事务收到半消息发送成功响应后Producer 执行自己本地数据库事务执行业务逻辑。本地事务执行结果三种状态1COMMIT_MESSAGE本地事务成功 → 通知 Broker 提交消息消息正式可见消费者可以消费2ROLLBACK_MESSAGE本地事务失败 → 通知 Broker 回滚消息直接删除永远不会被消费3UNKNOWN状态未知网络抖动、生产者宕机没返回 commit/rollbackBroker 需要事务回查3Broker 事务回查核心容错机制当 Broker 收到半消息长时间没有收到 commit/rollback 指令Broker 会定时主动回调 Producer 的checkLocalTransaction接口查询本地事务当前真实状态。1如果本地事务已经 commit → Broker 提交消息2如果本地事务已经 rollback → Broker 删除消息3如果还是 UNKNOWN等待下一次回查回查有最大次数超过直接回滚丢弃。⚠️关键点回查依赖业务实现检查本地事务状态接口如果生产者宕机重启重启后 Broker 依旧会来回查。【2】存储层面1半消息实际存储在 RMQ_SYS_TRANS_HALF_TOPIC 这个系统 topic不是业务 topic2commit 之后消息复制到真正业务 topic对外可见3rollback 直接删除半消息4事务回查 offset 存储在 RMQ_SYS_TRANS_OP_HALF_TOPIC。【3】时序图文字版Producer Broker|sendHalfMessage ------|保存半消息TRANSACTION_PREPARED不可消费|-------send ok --------||执行本地DB事务||----commit/rollback-----|更新消息状态||commit投递真实topicrollback删除# 如果生产者宕机没返回commit/rollbackBroker定时轮询半消息 Broker --回调checkLocalTransaction--Producer Producer查询数据库本地事务状态返回 Broker根据返回结果提交/回滚【四】优缺点✅优点1实现本地事务与消息发送原子性保证上游业务 DB 成功消息才对外可见2异步解耦下游业务异步执行提升吞吐量3有事务回查机制生产者宕机重启后依然可以完成事务状态确认容错能力强4基于 RocketMQ 原生支持不需要额外引入中间件。❌缺点1只保证发送端原子性消费端不保证。消费失败需要消费端重试 幂等2需要业务实现本地事务状态检查接口业务侵入业务必须提供查询本地事务执行状态的能力3事务回查会带来一定延迟不适合极致低延迟场景4如果业务没有做好幂等多次回查 消息重试会导致重复消费5不适合跨多个服务本地事务仅保证生产者本地 DB 发消息原子不能同时保证 A、B 两个服务 DB 同时原子6半消息会占用 broker 存储大量事务消息会增加 Broker 压力7回查频率、最大回查次数需要合理调参设置不当会出现消息悬停。常见坑业务没有实现 checkLocalTransaction或者查询逻辑写错导致消息一直回查最后丢弃业务消息丢失。【五】替换方案分布式事务方案对比最常用替代本地消息表方案很多生产环境不使用 RocketMQ 事务消息而是自己实现本地消息表避免 RocketMQ 事务消息的回查、半消息的各种坑。替代方案示例本地消息表很多项目不使用 RocketMQ 事务消息选择本地消息表。原理同一个本地事务同时写入业务数据 message_record 消息表然后定时任务扫描消息表发送 MQ更新消息状态。idbiz_id topic content status retry_count create_time1order_1001 ORDER_CREATE_TOPIC xxx0待发送02026‑08‑19status0 待发送1 发送成功2 发送失败TransactionalpublicvoidcreateOrderWithMsg(LongorderId){//1.本地事务内插入订单orderMapper.insert(order);//2.同一个事务插入消息记录状态待发送MessageRecordrecordnewMessageRecord();record.setBizId(orderId.toString());record.setTopic(ORDER_CREATE_TOPIC);record.setContent(xxx);record.setStatus(0);messageRecordMapper.insert(record);}定时任务扫描 status0重试次数少的消息发送 MQ发送成功更新 status1失败重试 retry_count。如果发送失败定时任务不断重试直到发送成功。优点不依赖 MQ 事务能力任意 MQ 都能用逻辑可控缺点需要维护一张消息表定时任务。【六】代码详细案例SpringBoot RocketMQ 事务消息【1】maven依赖dependencygroupIdorg.apache.rocketmq/groupIdartifactIdrocketmq-spring-boot-starter/artifactIdversion2.3.0/version/dependency【2】application.ymlrocketmq:name-server:127.0.0.1:9876producer:group:ORDER_TRANS_PRODUCER_GROUP业务场景创建订单本地数据库事务事务消息发送库存服务消费消息扣库存。数据库表t_order订单表order_id 作为事务唯一 id。【3】实现 RocketMQLocalTransactionListener两个核心方法executeLocalTransaction半消息发送成功后执行本地事务checkLocalTransactionBroker 回查检查本地事务状态importorg.apache.rocketmq.spring.annotation.RocketMQTransactionListener;importorg.apache.rocketmq.spring.core.RocketMQLocalTransactionListener;importorg.apache.rocketmq.spring.core.RocketMQLocalTransactionState;importorg.springframework.messaging.Message;importorg.springframework.stereotype.Component;RocketMQTransactionListener(txProducerGroupORDER_TRANS_PRODUCER_GROUP)ComponentpublicclassOrderTransactionListenerimplementsRocketMQLocalTransactionListener{// 注入自己的service操作数据库privatefinalOrderServiceorderService;publicOrderTransactionListener(OrderServiceorderService){this.orderServiceorderService;}/** * 半消息发送成功之后执行本地事务 * param msg mq消息 * param arg 业务自定义参数这里传入orderId * return */OverridepublicRocketMQLocalTransactionStateexecuteLocalTransaction(Messagemsg,Objectarg){LongorderId(Long)arg;try{// 执行本地数据库事务创建订单orderService.createOrder(orderId);// 本地事务成功通知Broker提交消息消息对消费者可见returnRocketMQLocalTransactionState.COMMIT;}catch(Exceptione){// 本地事务失败回滚消息删除returnRocketMQLocalTransactionState.ROLLBACK;}}/** * Broker事务回查生产者宕机没有返回commit/rollbackbroker回调这个接口查询本地事务状态 * 重点根据消息拿到业务唯一id查询数据库判断本地事务到底成功还是失败 */OverridepublicRocketMQLocalTransactionStatecheckLocalTransaction(Messagemsg){// 从消息header取出orderIdStringorderIdStr(String)msg.getHeaders().get(ORDER_ID);LongorderIdLong.parseLong(orderIdStr);// 查询数据库看订单是否创建成功booleanexistorderService.isOrderExist(orderId);if(exist){// DB已经成功提交消息returnRocketMQLocalTransactionState.COMMIT;}else{// DB没有数据事务回滚returnRocketMQLocalTransactionState.ROLLBACK;}}}【4】生产者发送事务消息importorg.apache.rocketmq.spring.core.RocketMQTemplate;importorg.springframework.messaging.Message;importorg.springframework.messaging.support.MessageBuilder;importorg.springframework.stereotype.Service;ServicepublicclassOrderMqService{privatefinalRocketMQTemplaterocketMQTemplate;publicstaticfinalStringORDER_TOPICORDER_CREATE_TOPIC;publicOrderMqService(RocketMQTemplaterocketMQTemplate){this.rocketMQTemplaterocketMQTemplate;}publicvoidsendOrderCreateTxMsg(LongorderId){MessageStringmessageMessageBuilder.withPayload(订单创建消息orderIdorderId).setHeader(ORDER_ID,orderId.toString()).build();//发送事务消息第二个参数是业务参数会传给executeLocalTransaction的argrocketMQTemplate.sendMessageInTransaction(ORDER_TOPIC,message,orderId);}}【5】OrderService 业务层ServicepublicclassOrderService{Transactional(rollbackForException.class)publicvoidcreateOrder(LongorderId){//本地事务插入订单OrderordernewOrder();order.setOrderId(orderId);order.setStatus(1);orderMapper.insert(order);// 模拟异常测试回滚// int i 1/0;}//给事务回查接口调用查询事务状态publicbooleanisOrderExist(LongorderId){returnorderMapper.selectById(orderId)!null;}}【6】消费者库存服务注意消费端必须做幂等⚠️重点RocketMQ 事务消息只保证发送端原子消费端会重试必须做幂等用 orderId 判断是否已经扣过库存。RocketMQMessageListener(topicORDER_CREATE_TOPIC,consumerGroupSTOCK_CONSUMER_GROUP)ComponentpublicclassStockConsumerimplementsRocketMQListenerString{privatefinalStockServicestockService;publicStockConsumer(StockServicestockService){this.stockServicestockService;}OverridepublicvoidonMessage(Stringmessage){//解析orderIdLongorderIdparseOrderId(message);//幂等判断如果该订单已经扣过库存直接return不再扣减if(stockService.hasDeduct(orderId)){return;}//执行扣库存stockService.deductStock(orderId);}privateLongparseOrderId(Stringmessage){//业务自行解析return1L;}}【七】注意事项1checkLocalTransaction回查接口必须保证幂等、查询数据库不能查缓存要查真实 DB 状态2消费端一定要做幂等事务消息解决发送端原子消费失败会重试3合理设置 Broker 事务回查参数transactionTimeOut、checkMaxTimes4不要把业务耗时很长的逻辑放在executeLocalTransaction半消息超时会触发大量回查5事务 producer group 必须唯一一个 group 对应一组事务监听6如果回查多次还是 UNKNOWNbroker 会直接回滚丢弃消息业务要监控事务消息丢弃告警。常见问题QRocketMQ 事务消息能保证消费一定成功吗A不能。只能保证消息一定投递到 broker上游 DB 成功消息才可见消费失败会重试消费成功需要业务自己保证。Q半消息会占用磁盘吗A会存储在系统半消息 topiccommit/rollback 之后才清理。Q和本地消息表怎么选A简单业务可以用 RocketMQ 事务消息复杂生产系统追求可控性优先本地消息表。