RocketMQ 事务消息概述:为什么需要事务消息
RocketMQ 事务消息概述:为什么需要事务消息
在分布式系统中,一个经典难题是:本地事务执行和消息发送是两个独立操作,无法保证原子性。先执行本地事务再发送消息,如果消息发送失败,下游系统无法感知数据变更。先发送消息再执行本地事务,如果本地事务回滚,已发送的消息无法撤回。
RocketMQ 的事务消息机制为解决这个问题提供了标准方案。它保证了本地事务执行与消息发送的原子性,即二者要么同时成功,要么同时失败。
---
2. 事务消息的执行流程:半消息、本地事务与回查
RocketMQ 事务消息分为三个核心阶段:发送半消息、执行本地事务、事务状态回查。
第一阶段:发送半消息
生产者向 Broker 发送一条半消息。半消息与普通消息的区别在于,它此时对消费者不可见。Broker 将消息存储后,返回发送成功,但消息的状态标记为“待确认”。消费者无论使用 push 还是 pull 模式,都无法获取到此消息。
第二阶段:执行本地事务
生产者收到半消息发送成功的响应后,执行本地事务逻辑。根据本地事务的执行结果,生产者向 Broker 发送二次确认:提交或回滚。
- 本地事务成功:发送 COMMIT 指令,Broker 将半消息标记为可消费,消费者可以拉取到该消息。
- 本地事务失败:发送 ROLLBACK 指令,Broker 删除该半消息,消费者永远看不到这条消息。
第三阶段:事务状态回查
如果生产者因崩溃、网络超时等故障,未能向 Broker 发送二次确认,Broker 会主动向生产者发起回查。回查的目的是让生产者再次确认该消息对应的事务最终状态。生产者需要实现回查接口,根据本地事务的执行记录返回 COMMIT 或 ROLLBACK。
---
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 根据回查结果删除了半消息。
- 实际上本地事务稍后提交成功,但消息已被删除,消费者永远不会收到通知。
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) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 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 RocketMQListener<String> { @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.group=order-tx-producer-group rocketmq.producer.send-message-timeout=5000 rocketmq.producer.retry-times-when-send-failed=3 # 事务消息回查配置 rocketmq.producer.check-request-hold-max=2000 rocketmq.producer.transaction-check-interval=60000send-message-timeout控制半消息发送的超时时间,建议 3 到 5 秒。transaction-check-interval控制 Broker 发起回查的最小间隔,默认 60 秒。对于实时性要求高的业务,可以适当调低此值,但需注意这会增加回查频率和系统负载。
6.4 监控与告警
重点监控以下指标:
- 半消息数量:如果持续增长且不下降,说明大量事务消息处于待确认状态,可能存在生产者故障。
- 回查次数:频繁回查表明生产者响应不及时或返回 UNKNOW 比例过高。
- 事务状态表中状态为 0 的记录数量:如果长时间存在大量状态为 0 的记录,说明部分本地事务执行耗时过长或发生了死锁。
- 消费者端的消费失败率:事务消息的重复投递可能导致消费端压力增大。
通过事务状态记录表和 Broker 的回查机制配合,RocketMQ 事务消息在绝大多数故障场景下都能保证本地事务和消息发送的最终一致性。理解回查的触发时机和事务状态表的角色,是排查“本地事务成功但消息未发送”问题的关键所在。
