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 事务消息在绝大多数故障场景下都能保证本地事务和消息发送的最终一致性。理解回查的触发时机和事务状态表的角色是排查“本地事务成功但消息未发送”问题的关键所在。

相关新闻

高通FAS调度技术:提升移动设备性能与能效

高通FAS调度技术:提升移动设备性能与能效

1. 高通处理器FAS调度技术解析最近在移动设备性能优化领域,FAS调度技术引起了广泛关注。作为一名长期从事移动设备性能调优的工程师,我发现这套调度机制在高通平台上展现出惊人的潜力——既能提升游戏帧率,又能降低日常使用功耗。不同于传统的…

2026/9/22 7:10:10 阅读更多 →
2026 AI Agent落地实操:5个真跑通场景,省人省时还能复用!

2026 AI Agent落地实操:5个真跑通场景,省人省时还能复用!

不是聊趋势,是看哪些业务已经开始稳定省人、省时、能复用。一个很现实的反差:很多团队还在写“AI战略PPT”,但另一批团队已经让 Agent 接手周报、客服、销售线索、知识检索和内容分发。差距不在模型能力,而在有没有把任务拆成可执…

2026/9/24 0:01:14 阅读更多 →
Kubernetes 多集群管理:AI 平台跨集群调度的工程实践

Kubernetes 多集群管理:AI 平台跨集群调度的工程实践

Kubernetes 多集群管理:AI 平台跨集群调度的工程实践 一、单集群不是银弹:为什么 AI 平台必须走向多集群 任何一个从零开始搭建的 AI 推理平台,最初都会把全部服务部署在一个 Kubernetes 集群里。这个阶段集群规模可控——几个推理服务、一个…

2026/9/22 21:17:51 阅读更多 →

最新新闻

工业无线通信模块哪个品牌靠谱?2026 选型避坑指南

工业无线通信模块哪个品牌靠谱?2026 选型避坑指南

核心结论工业无线通信模块选型,品牌可靠性需从抗干扰、稳定性、认证资质、服务支持四个维度综合评估,而非只看参数表。据中国信通院数据,2026 年国内物联网通信市场规模预计约 1.87 万亿元,工业物联网贡献约 1.98 万亿元。工程师常…

2026/9/24 8:47:01 阅读更多 →
EMC检测费用七层动态模型解析

EMC检测费用七层动态模型解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/24 8:47:01 阅读更多 →
4针风扇PWM调速电路设计:从原理到PCB实战

4针风扇PWM调速电路设计:从原理到PCB实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/24 8:47:01 阅读更多 →
【Pi 超轻量插件】让 Jev 去掉无用的工具调用结果,留一个清爽的上下文!

【Pi 超轻量插件】让 Jev 去掉无用的工具调用结果,留一个清爽的上下文!

背景与痛点 长会话里那些“食之无味”的历史工具输出 在使用 AI 进行较复杂的重构或排查任务时,都会遇到上下文过多压缩的情况。 在 Pi 中虽然说有压缩命令,但我们可以看看他具体压缩了些什么。 如果你用 Pi 跑过时间稍长的任务,大概清楚那个…

2026/9/24 8:47:01 阅读更多 →
PaddleNLP T5 模型家族全解析:预训练权重清单与源码级使用指南

PaddleNLP T5 模型家族全解析:预训练权重清单与源码级使用指南

人工智能大模型预训练微调LoRARLHF强化学习分布式训练 【免费下载链接】PaddleNLP Easy-to-use and powerful LLM and SLM library with awesome model zoo. 项目地址: https://gitcode.com/gh_mirrors/pa/PaddleNLP 点击查看 免费下载 T5(Text-to-Text…

2026/9/24 8:47:01 阅读更多 →
DiT详解

DiT详解

前言 在Sora[1]的技术报告中,作者指出Sora是一个Diffusion Transformer。这个Diffusion Transformer便是我们这里将要介绍的DiT[2]。相较于我们之前介绍的LDM[3],DiTs也是作用在潜空间,它最大的改进是将U-Net的CNN替换为了Transformer。同时…

2026/9/24 8:46:01 阅读更多 →

日新闻

基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程

基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程

简介:这是一套面向计算机、人工智能、自动化等专业学生与教师的毕业设计级项目资源,围绕YOLOv8实现渔船作业监控系统,可用于毕设、课程设计、大作业或项目立项演示。压缩包共97个文件,约24.21MB,以70个Python源码文件为…

2026/9/24 0:00:19 阅读更多 →
单细胞注释实战:基于Scanpy的标记基因与参考映射流程解析

单细胞注释实战:基于Scanpy的标记基因与参考映射流程解析

简介:一份基于单细胞RNA测序数据的细胞类型注释算法研究Python毕业设计源码,针对计算机相关专业正在做毕设或需要项目实战的学习者,可用于课程设计与期末大作业。项目代码完整、经导师指导评审通过,可直接运行,覆盖数据…

2026/9/24 0:00:19 阅读更多 →
C#源生成器实战:用增量生成器替代反射,告别AOT崩溃

C#源生成器实战:用增量生成器替代反射,告别AOT崩溃

第一次在项目里被反射卡住,是在一个老旧的WinForms模块里:几十个类依赖PropertyChanged通知,运行时反射读属性、发通知,每次启动慢半拍不说,一上.NET Native/AOT裁剪模式几乎全面崩盘。后来我把这段逻辑全部改成C#源生…

2026/9/24 0:00:19 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/23 4:55:02 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/23 4:49:06 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/23 9:53:41 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/23 9:53:40 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/23 9:53:40 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/23 9:53:40 阅读更多 →