当前位置: 首页 > news >正文

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内部触发了事务回查,而回查逻辑需要查询数据库中的订单记录,就会产生竞态条件:


  1. 半消息发送成功,Broker 存储了待确认消息。
  2. 生产者因网络延迟、GC 停顿或进程崩溃,未能在超时前发送 COMMIT。
  3. Broker 触发回查,生产者查询本地数据库确认事务状态。
  4. 如果此时本地事务尚未提交,数据库中没有订单记录,回查返回 UNKNOW 或 ROLLBACK。
  5. Broker 根据回查结果删除了半消息。
  6. 实际上本地事务稍后提交成功,但消息已被删除,消费者永远不会收到通知。


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) ) 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=60000


send-message-timeout控制半消息发送的超时时间,建议 3 到 5 秒。transaction-check-interval控制 Broker 发起回查的最小间隔,默认 60 秒。对于实时性要求高的业务,可以适当调低此值,但需注意这会增加回查频率和系统负载。


6.4 监控与告警


重点监控以下指标:

  • 半消息数量:如果持续增长且不下降,说明大量事务消息处于待确认状态,可能存在生产者故障。
  • 回查次数:频繁回查表明生产者响应不及时或返回 UNKNOW 比例过高。
  • 事务状态表中状态为 0 的记录数量:如果长时间存在大量状态为 0 的记录,说明部分本地事务执行耗时过长或发生了死锁。
  • 消费者端的消费失败率:事务消息的重复投递可能导致消费端压力增大。


通过事务状态记录表和 Broker 的回查机制配合,RocketMQ 事务消息在绝大多数故障场景下都能保证本地事务和消息发送的最终一致性。理解回查的触发时机和事务状态表的角色,是排查“本地事务成功但消息未发送”问题的关键所在。

http://www.jsqmd.com/news/1247714/

相关文章:

  • ShardingSphere与Seata AT分布式事务整合实践
  • DRV8353Rx-EVM GUI 实战指南:从零上手磁场定向控制(FOC)电机驱动
  • 程序员转型AI的四大路径与6个月速成指南
  • AI模型推理参数调优:精度、速度与显存的平衡艺术
  • OES支F协议解析:Web3开发的核心标准与实践
  • 盘点沈阳浑南画室靠谱供应商推荐好物榜 - 招财兔数字员工
  • 深度学习模型优化:稀疏计算与结构化剪枝实践
  • 《从0到1搭建私域社群SOP手册》免费领:一个人也能搭出高转化社群
  • OpenGL开发环境搭建教程
  • AI 原型生成:从 PRD 到可交互前端的自动化验证闭环
  • 收到的PDF加了密码,输入密码才能打开?其实有两道锁,很多人只遇到第一道
  • 弱电视频监控系统技术要求与架构设计全解析
  • 三易串口屏VP开发环境深度解析:为什么会C语言,就能快速开发工业HMI
  • UE5与WEB双向通信实战:基于WebUI插件实现数据可视化与交互
  • AI学术工具助力研究生高效论文写作
  • 中小企业ERP系统对接实战:以管家婆进销存为例(附Checklist)
  • 新手第一次在怀化卖黄金必读!避开回收套路,6 家正规门店汇总(2026 年 7 月更新) - 不晚生活号
  • AI技能开发实战:从僵尸文件到效率神器的五大标准
  • 计算机Python毕设实战-网络音乐资源播放与歌单管理平台 基于 Python Web 的智能音乐娱乐平台【完整源码+LW+部署说明+演示视频,全bao一条龙等】
  • 亲身探访长沙卡地亚售后服务中心|最新电话及地址(2026年7月最新) - 卡地亚服务中心
  • 2026年7月南宁劳力士回收价格查询:我找了5家商家,哪个渠道收的价格更高?客户反馈+实测攻略全公开! - 嘉价奢侈品回收平台
  • AI智能体核心架构与行业应用解析
  • 医药AIGC实战:AI疾病筛查技术解析与应用
  • AI Agent架构设计与核心组件解析
  • Kimi K3登顶前端代码竞技场:AI编程助手的实战应用指南
  • Linux 实时优化:禁用内核调试 / 跟踪功能实战教程
  • 抖音小店一件代发需要准备哪些工具? - 抖掌柜
  • BQ41Z50数据闪存参数详解:Gas Gauging与RA Table配置实战
  • 关于脉冲电流源中开关mos管产生振铃现象的分析
  • 微信消息撤回机制解析与使用技巧