RocketMQ 事务消息概述:为什么需要事务消息 RocketMQ 事务消息概述为什么需要事务消息在分布式系统中一个经典难题是本地事务执行和消息发送是两个独立操作无法保证原子性。先执行本地事务再发送消息如果消息发送失败下游系统无法感知数据变更。先发送消息再执行本地事务如果本地事务回滚已发送的消息无法撤回。RocketMQ 的事务消息机制为解决这个问题提供了标准方案。它保证了本地事务执行与消息发送的原子性即二者要么同时成功要么同时失败。RocketMQ事务消息半消息机制 本地事务 回查本地事务和消息发送原子性保障, 最终一致无事务消息的困境1. 先执行本地事务, 后发消息本地事务成功 消息发送失败 数据不一致2. 先发消息, 后执行本地事务消息发送成功 本地事务回滚 数据不一致---2. 事务消息的执行流程半消息、本地事务与回查RocketMQ 事务消息分为三个核心阶段发送半消息、执行本地事务、事务状态回查。第一阶段发送半消息生产者向 Broker 发送一条半消息。半消息与普通消息的区别在于它此时对消费者不可见。Broker 将消息存储后返回发送成功但消息的状态标记为“待确认”。消费者无论使用 push 还是 pull 模式都无法获取到此消息。第二阶段执行本地事务生产者收到半消息发送成功的响应后执行本地事务逻辑。根据本地事务的执行结果生产者向 Broker 发送二次确认提交或回滚。本地事务成功发送 COMMIT 指令Broker 将半消息标记为可消费消费者可以拉取到该消息。本地事务失败发送 ROLLBACK 指令Broker 删除该半消息消费者永远看不到这条消息。第三阶段事务状态回查如果生产者因崩溃、网络超时等故障未能向 Broker 发送二次确认Broker 会主动向生产者发起回查。回查的目的是让生产者再次确认该消息对应的事务最终状态。生产者需要实现回查接口根据本地事务的执行记录返回 COMMIT 或 ROLLBACK。消费者RocketMQ Broker生产者消费者RocketMQ Broker生产者第一阶段发送半消息消费者此时不可见此消息第二阶段执行本地事务消费者永远不会看到此消息等待超时触发回查alt[本地事务执行成功][本地事务执行失败][生产者故障未发送确认]1. 发送半消息2. 存储消息标记为待确认3. 半消息发送成功4. 执行本地事务5a. 发送 COMMIT6a. 标记消息为可消费7a. 拉取消息并消费5b. 发送 ROLLBACK6b. 删除半消息5c. 发起事务回查6c. 返回本地事务执行结果7c. 根据结果提交或回滚---3. 本地事务成功但消息未发送的根因分析生产环境中最常见的故障场景是数据库数据已经变更但消费者迟迟没有收到消息。问题根源在于事务消息的实现中存在一个关键的时序陷阱。3.1 典型的问题代码Transactional public void processOrder(Order order) { // 1. 发送半消息 TransactionSendResult result rocketMQTemplate.sendMessageInTransaction( order-tx-group, order-topic, message, order); // 2. 执行业务逻辑 orderRepository.save(order); // 3. 本地事务方法返回框架根据返回结果决定 COMMIT 或 ROLLBACK }这段代码的问题在于sendMessageInTransaction方法在Transactional注解修饰的方法内部被调用。如果sendMessageInTransaction内部触发了事务回查而回查逻辑需要查询数据库中的订单记录就会产生竞态条件半消息发送成功Broker 存储了待确认消息。生产者因网络延迟、GC 停顿或进程崩溃未能在超时前发送 COMMIT。Broker 触发回查生产者查询本地数据库确认事务状态。如果此时本地事务尚未提交数据库中没有订单记录回查返回 UNKNOW 或 ROLLBACK。Broker 根据回查结果删除了半消息。实际上本地事务稍后提交成功但消息已被删除消费者永远不会收到通知。RocketMQ Broker本地数据库生产者RocketMQ Broker本地数据库生产者数据在事务中尚未提交, 对外不可见网络抖动或 GC 停顿未能及时发送 COMMIT超时未收到确认触发回查数据已持久化但消息已被删除, 下游无法感知1. 发送半消息半消息发送成功2. 开始本地事务执行 INSERT 订单3. 回查事务状态4. 查询订单记录5. 未查询到记录本地事务尚未提交6. 返回 ROLLBACK7. 删除半消息8. 本地事务提交成功3.2 问题的本质事务消息的回查机制与本地事务的提交时机之间存在时间窗口。在这个窗口内回查逻辑无法正确判断本地事务的最终状态。解决这个问题的关键在于让回查逻辑能够在本地事务提交之前就能准确预判事务的最终结果。---4. 生产级解决方案事务状态记录表模式4.1 核心设计思路在本地事务内部与业务数据一同写入一条事务状态记录。这条记录与业务数据同在一个本地事务中要么一起提交成功要么一起回滚。回查时通过查询事务状态记录来判定本地事务的执行结果而不是依赖业务数据是否存在。这样做的好处是业务数据可能尚未提交而对回查不可见但通过Transactional保护的事务状态记录要么已经提交要么在事务回滚时一同被清除。回查逻辑看到一个确定的、一致的状态。4.2 事务状态记录表结构CREATE TABLE t_transaction_log ( transaction_id VARCHAR(64) NOT NULL PRIMARY KEY COMMENT 事务 ID关联消息的 transactionId, status TINYINT NOT NULL DEFAULT 0 COMMENT 事务状态: 0-执行中, 1-已提交, 2-已回滚, business_key VARCHAR(128) COMMENT 业务主键如订单号, create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, INDEX idx_business_key (business_key), INDEX idx_create_time (create_time) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT分布式事务状态记录表;4.3 完整的生产者端实现本地事务执行监听器。这个类是事务消息的核心负责执行本地事务并通知 Broker 最终结果Slf4j Component public class OrderTransactionListener implements RocketMQLocalTransactionListener { Autowired private OrderService orderService; Autowired private TransactionLogService transactionLogService; Override public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) { String transactionId msg.getTransactionId(); String orderJson new String((byte[]) msg.getBody(), StandardCharsets.UTF_8); Order order JSON.parseObject(orderJson, Order.class); try { // 执行本地事务业务数据与事务日志在同一个本地事务中写入 orderService.createOrderWithTransactionLog(transactionId, order); log.info(本地事务执行成功, transactionId: {}, transactionId); return RocketMQLocalTransactionState.COMMIT; } catch (Exception e) { log.error(本地事务执行失败, transactionId: {}, transactionId, e); return RocketMQLocalTransactionState.ROLLBACK; } } Override public RocketMQLocalTransactionState checkLocalTransaction(Message msg) { String transactionId msg.getTransactionId(); // 回查时查询事务状态记录表 TransactionLog log transactionLogService.getByTransactionId(transactionId); if (log null) { log.warn(回查未找到事务记录, transactionId: {}, 返回 UNKNOW, transactionId); return RocketMQLocalTransactionState.UNKNOW; } if (log.getStatus() 1) { log.info(回查确认事务已提交, transactionId: {}, transactionId); return RocketMQLocalTransactionState.COMMIT; } if (log.getStatus() 2) { log.info(回查确认事务已回滚, transactionId: {}, transactionId); return RocketMQLocalTransactionState.ROLLBACK; } log.warn(回查事务状态未知, transactionId: {}, status: {}, transactionId, log.getStatus()); return RocketMQLocalTransactionState.UNKNOW; } }业务服务层。关键点在于Transactional保证了事务日志和业务数据在同一本地事务中Slf4j Service public class OrderService { Autowired private OrderMapper orderMapper; Autowired private TransactionLogMapper transactionLogMapper; Transactional(rollbackFor Exception.class) public void createOrderWithTransactionLog(String transactionId, Order order) { // 1. 插入事务状态记录状态为 0 (执行中) TransactionLog txLog new TransactionLog(); txLog.setTransactionId(transactionId); txLog.setStatus(0); txLog.setBusinessKey(order.getOrderNo()); transactionLogMapper.insert(txLog); // 2. 执行业务逻辑 orderMapper.insert(order); // 3. 业务执行成功后更新事务状态为 1 (已提交) transactionLogMapper.updateStatus(transactionId, 1); log.info(本地事务执行完成, transactionId: {}, orderNo: {}, transactionId, order.getOrderNo()); } }事务日志服务Service public class TransactionLogService { Autowired private TransactionLogMapper transactionLogMapper; public TransactionLog getByTransactionId(String transactionId) { return transactionLogMapper.selectByTransactionId(transactionId); } }消息发送入口Service public class OrderMessageService { Autowired private RocketMQTemplate rocketMQTemplate; public void sendOrderMessage(Order order) { String transactionId UUID.randomUUID().toString(); MessageBuilder builder MessageBuilder.withPayload(JSON.toJSONBytes(order)); builder.setHeader(RocketMQHeaders.TRANSACTION_ID, transactionId); rocketMQTemplate.sendMessageInTransaction( order-tx-producer-group, order-topic, builder.build(), order ); log.info(事务消息发送请求已提交, transactionId: {}, orderNo: {}, transactionId, order.getOrderNo()); } }生产者配置Configuration public class RocketMQConfig { Value(${rocketmq.name-server}) private String nameServer; Bean public RocketMQTemplate rocketMQTemplate() { RocketMQTemplate template new RocketMQTemplate(); template.setNameServer(nameServer); return template; } }4.4 消费者端实现消费者端需要对事务消息做幂等处理因为回查可能导致消息被重复投递Slf4j Service RocketMQMessageListener( topic order-topic, consumerGroup order-consumer-group, selectorExpression *) public class OrderMessageConsumer implements RocketMQListenerString { Autowired private InventoryService inventoryService; Override public void onMessage(String message) { Order order JSON.parseObject(message, Order.class); // 幂等性检查根据 orderNo 判断是否已处理过 if (inventoryService.isAlreadyProcessed(order.getOrderNo())) { log.info(订单已处理跳过重复消息, orderNo: {}, order.getOrderNo()); return; } try { inventoryService.deductStock(order); log.info(订单消费成功, orderNo: {}, order.getOrderNo()); } catch (Exception e) { log.error(订单消费失败, orderNo: {}, order.getOrderNo(), e); // 抛出异常触发 RocketMQ 的重试机制 throw new RuntimeException(消费失败, e); } } }---5. 事务消息回查机制的深入理解5.1 回查的触发条件RocketMQ Broker 在以下情况会触发事务回查半消息存储后在指定时间内未收到生产者的二次确认。默认超时时间为 6 秒。生产者返回了 UNKNOW 状态Broker 会定期回查直到获得明确的 COMMIT 或 ROLLBACK。5.2 回查的频率控制Broker 对同一条消息的回查不是无限次的。默认配置下回查间隔从 60 秒开始如果生产者持续返回 UNKNOW后续回查的间隔逐步放大。最大回查次数默认 15 次。超过最大次数后如果仍无法得到确定结果Broker 会将该消息丢弃。这个机制避免了因代码逻辑错误导致的无限制回查。5.3 回查中的线程安全回查请求可能由 Broker 端的多个线程并发发出。生产者的回查处理逻辑必须是线程安全的。这要求事务状态记录表的查询和更新操作具备原子性。使用数据库的唯一索引和行锁可以天然保证这一点。---6. 生产环境注意事项6.1 事务状态记录的定期清理随着业务运行t_transaction_log表会持续增长。需要定期清理已完成的事务记录避免表过大影响查询性能-- 清理 7 天前已完成的记录 DELETE FROM t_transaction_log WHERE status IN (1, 2) AND create_time DATE_SUB(NOW(), INTERVAL 7 DAY) LIMIT 10000;将这条 SQL 配置为定时任务在业务低峰期分批执行避免长事务锁表。6.2 消费者幂等性保障事务消息的回查和重试机制意味着消息可能被投递多次。消费者必须实现幂等性基于业务主键如订单号进行去重。使用 Redis 缓存已处理的消息 ID设置过期时间与业务窗口匹配。数据库唯一约束兜底确保即使去重逻辑失效数据也不会重复写入。6.3 超时与重试配置建议# Producer 配置 rocketmq.producer.grouporder-tx-producer-group rocketmq.producer.send-message-timeout5000 rocketmq.producer.retry-times-when-send-failed3 # 事务消息回查配置 rocketmq.producer.check-request-hold-max2000 rocketmq.producer.transaction-check-interval60000send-message-timeout控制半消息发送的超时时间建议 3 到 5 秒。transaction-check-interval控制 Broker 发起回查的最小间隔默认 60 秒。对于实时性要求高的业务可以适当调低此值但需注意这会增加回查频率和系统负载。6.4 监控与告警重点监控以下指标半消息数量如果持续增长且不下降说明大量事务消息处于待确认状态可能存在生产者故障。回查次数频繁回查表明生产者响应不及时或返回 UNKNOW 比例过高。事务状态表中状态为 0 的记录数量如果长时间存在大量状态为 0 的记录说明部分本地事务执行耗时过长或发生了死锁。消费者端的消费失败率事务消息的重复投递可能导致消费端压力增大。通过事务状态记录表和 Broker 的回查机制配合RocketMQ 事务消息在绝大多数故障场景下都能保证本地事务和消息发送的最终一致性。理解回查的触发时机和事务状态表的角色是排查“本地事务成功但消息未发送”问题的关键所在。

本月热点