资讯详情

资讯详情

RocketMQ 事务消息:订单与积分的最终一致

RocketMQ 事务消息订单与积分的最终一致作者鱼宵 实战驱动系列 · 第 6 篇完整课程与可运行源码已开源在 Giteehttps://gitee.com/j67mk2/rocketmq-journey 本文对应 lesson-06/你有没有碰到过这种线上事故——订单库明明写进去了积分系统却一直没加上分别急着背锅这不是谁的代码写错了是跨系统一致性的经典难题。lesson-06 工程clone 下来就有里就有两条消息一条正常提交、一条走回查跑一遍你就能亲眼看到半消息从隐身到现身的全过程。一、下单成功了积分却没加上——这锅谁背你写过下单接口吧用户点支付成功你的系统要干两件事往订单库插一条订单记录本地数据库事务发一条 MQ 消息通知积分系统给这个用户加 100 积分。这两件事分属两个系统没法用同一个数据库事务包起来。于是两种错误姿势先发消息、再写库。消息发出去了积分系统都把积分加了结果写库失败回滚了——积分凭空多了订单却不存在用户白嫖。先写库、再发消息。订单写进去了正要发消息时机器宕机重启——订单有了积分永远没加用户来投诉。这就是分布式环境下的原子性问题两个系统的两个操作怎么保证结果一致别慌RocketMQ 专门给你准备了一把钥匙事务消息。它的核心思路说出来你就懂——跟淘宝的担保交易一模一样。二、担保交易类比半消息到底是什么你在淘宝下单钱不是直接打给卖家而是先冻结在支付宝里。卖家真正发货了支付宝才把钱打给他卖家一直不发货支付宝会退钱。RocketMQ 的事务消息就是这个逻辑概念担保交易类比人话解释Half 半消息付了定金但冻结在支付宝里的钱已写到 Broker但存在内部 Half Topic 里消费者完全看不到本地事务卖家真正发货你自己业务库里的操作插订单、扣库存COMMIT确认打款给卖家消息正式投递到真实 Topic消费者可见ROLLBACK交易取消、退定金半消息丢弃消费者永远看不到UNKNOWN卖家说我不确定发没发货不提交也不回滚先挂着等回查回查 checkLocalTransaction支付宝上门问卖家货到底发了没Broker 定时回调生产者那条消息的本地事务到底成了没一句话半消息 冻结的定金。消费者看不到它直到你说确认成交COMMIT。三、两阶段流程从发消息到消费中间发生了什么Producer你的服务 BrokerMQ Consumer积分系统 │ │ │ │ ① 发 Half 半消息 │ │ │ ───────────────────────────▶ │ 落盘到内部 Half 队列 │ │ │ 对消费者隐身 │ │ ② Broker 确认收到了 │ │ │ ◀─────────────────────────── │ │ │ ③ 执行本地事务 │ │ │ 写订单表…… │ │ │ ④ 向 Broker 报告结果 │ │ │ ───────────────────────────▶ │ COMMIT投递到真实 Topic ─▶ │ ⑤ 消费到消息 │ │ ROLLBACK丢弃半消息 │ 永远收不到 │ │ UNKNOWN先挂着等回查 │ │ │ │ │ ⑥ 超时后 Broker 主动回查 │ │ │ ◀─────────────────────────── │ 每隔 10 秒上门一次 │ │ ⑦ 查库后回答 COMMIT/ROLLBACK │ │ │ ───────────────────────────▶ │ COMMIT投递 ──────────────▶ │ ⑧ 这时才收到第 ③ 步执行本地事务时可能因为数据库超时、应用宕机等原因你当时拿不准本地事务到底成没成。这时返回 UNKNOWN。Broker 也不敢替你做决定——它不知道你的业务库最后是什么样。于是每隔transactionCheckInterval本教程特意配成了10 秒就是为了演示Broker 会回调生产者的checkLocalTransaction方法问一句“那条消息的本地事务最后成了没”真实项目里这个方法干的事就是查数据库订单表里有这条记录吗有 → COMMIT没有 → ROLLBACK。回查有次数上限默认 15 次一直没答案就丢弃。注意回查是 Broker主动回调生产者 JVM的。所以 Demo 生产者发完消息后不能立刻退出得保持存活等回查上门——这是本课最容易踩的坑。四、动手完整代码跑通两条路径环境前提RocketMQ 三件套已启动NameServer Broker 控制台工程已mvn clean install编译过。第 1 步先启动消费者重要顺序为什么先起消费者新消费组默认只收启动之后产生的新消息。而且消息 B 要等回查才会投递消费者必须一直在场等着。[Console]::OutputEncoding[System.Text.Encoding]::UTF8$env:JAVA_TOOL_OPTIONS-Dfile.encodingUTF-8$env:JAVA_HOMEC:\Program Files\Java\jdk-17cd rocketmq-journey\lesson-06 mvn exec:java-Dexec.mainClasscom.example.lesson06.TransactionConsumer看到消费者已启动正在等待消息……就保持这个窗口别关。第 2 步再启动生产者新开一个 PowerShell 窗口同样设好环境变量后执行mvn exec:java-Dexec.mainClasscom.example.lesson06.TransactionProducer生产者会保持存活约 30 秒等回查看到收到 Broker 回查后才退出。这期间切回消费者窗口能看到两条消息先后到达。TransactionProducer.java —— 事务消息生产者完整代码packagecom.example.lesson06;importorg.apache.rocketmq.client.producer.LocalTransactionState;importorg.apache.rocketmq.client.producer.TransactionListener;importorg.apache.rocketmq.client.producer.TransactionMQProducer;importorg.apache.rocketmq.client.producer.TransactionSendResult;importorg.apache.rocketmq.common.message.Message;importorg.apache.rocketmq.common.message.MessageExt;importjava.nio.charset.StandardCharsets;importjava.text.SimpleDateFormat;importjava.util.Date;/** * 第 6 课 Demo事务消息生产者 * * 干的事往 TopicLesson06 发 2 条事务消息演示两条路径 * - 消息AtagTagCommit本地事务一次就成功 → 返回 COMMIT * 消息立刻对消费者可见正常提交路径 * - 消息BtagTagCheckback本地事务结果未知模拟超时/宕机→ 返回 UNKNOWN * 约 10 秒后 Broker 主动回查本 JVM我们在回查里判定成功并 COMMIT回查路径。 */publicclassTransactionProducer{publicstaticvoidmain(String[]args)throwsException{SimpleDateFormatsdfnewSimpleDateFormat(HH:mm:ss.SSS);// 1. 创建【事务消息专用】生产者 TransactionMQProducer。// 它和普通 DefaultMQProducer 的区别多挂一个事务监听器// Broker 回查时就是回调这个监听器里的方法。TransactionMQProducerproducernewTransactionMQProducer(lesson06_producer_group);// 2. 告诉生产者 NameServer 在哪NameServer 电话簿先问路再连 Brokerproducer.setNamesrvAddr(127.0.0.1:9876);// 3. 注册事务监听器两个方法会在两个时机被回调producer.setTransactionListener(newTransactionListener(){/** 半消息落盘后立刻回调在这里执行本地事务并报告结果 */OverridepublicLocalTransactionStateexecuteLocalTransaction(Messagemsg,Objectarg){StringbodynewString(msg.getBody(),StandardCharsets.UTF_8);if(TagCommit.equals(msg.getTags())){// 路径一本地事务一次成功 → COMMIT消息正式投递给消费者System.out.println(sdf.format(newDate()) [本地事务] 消息A 执行成功模拟订单已落库返回 COMMIT);returnLocalTransactionState.COMMIT_MESSAGE;}else{// 路径二模拟拿不准结果超时/宕机→ UNKNOWN等回查System.out.println(sdf.format(newDate()) [本地事务] 消息B 结果未知模拟本地操作超时返回 UNKNOWN等待 Broker 回查);returnLocalTransactionState.UNKNOW;}}/** UNKNOWN 之后Broker 每隔 10 秒上门回查时回调这里 */OverridepublicLocalTransactionStatecheckLocalTransaction(MessageExtmsg){StringbodynewString(msg.getBody(),StandardCharsets.UTF_8);// 教程简化回查一律判定成功真实代码这里是查订单表有没有这条记录System.out.println(sdf.format(newDate()) [回查] 收到 Broker 回查消息body → 查库后判定成功返回 COMMIT);returnLocalTransactionState.COMMIT_MESSAGE;}});// 4. 启动生产者producer.start();System.out.println(sdf.format(newDate()) 事务生产者已启动。);// 5. 发消息A注意用的是 sendMessageInTransaction不是普通的 sendMessagemsgAnewMessage(TopicLesson06,TagCommit,消息A用户下单成功请给用户加 100 积分.getBytes(StandardCharsets.UTF_8));TransactionSendResultresultAproducer.sendMessageInTransaction(msgA,null);System.out.println(sdf.format(newDate()) 发送消息A 完成resultA);// 6. 发消息B走 UNKNOWN → 回查路径MessagemsgBnewMessage(TopicLesson06,TagCheckback,消息B用户下单成功请给用户发短信通知.getBytes(StandardCharsets.UTF_8));TransactionSendResultresultBproducer.sendMessageInTransaction(msgB,null);System.out.println(sdf.format(newDate()) 发送消息B 完成resultB);// 7. 重点回查是 Broker 回调【本 JVM】的。现在退出的话// 10 秒后 Broker 找不到生产者消息B 永远卡在半消息状态。// 保持存活 30 秒等回查上门。System.out.println(生产者发送完毕保持 JVM 存活 30 秒等待 Broker 回查消息B……);for(inti1;i6;i){Thread.sleep(5000);System.out.println(已等待 (i*5) 秒……);}producer.shutdown();System.out.println(sdf.format(newDate()) 生产者已关闭。);}}这段代码的三个关键点TransactionMQProducer事务消息专用生产者。普通DefaultMQProducer没有事务监听器这个挂载点。sendMessageInTransaction(msg, null)事务消息的发送入口。它内部先把 Half 半消息发给 Broker等 Broker 确认后再回调executeLocalTransaction最后根据返回值决定 COMMIT / ROLLBACK / 挂起。LocalTransactionState.UNKNOW注意拼写是 UNKNOW少一个 N这是 RocketMQ API 里的历史拼写照抄即可。TransactionConsumer.java —— 事务消息消费者完整代码packagecom.example.lesson06;importorg.apache.rocketmq.client.consumer.DefaultMQPushConsumer;importorg.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;importorg.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;importorg.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;importorg.apache.rocketmq.common.message.MessageExt;importjava.nio.charset.StandardCharsets;importjava.text.SimpleDateFormat;importjava.util.Date;importjava.util.List;importjava.util.concurrent.atomic.AtomicInteger;publicclassTransactionConsumer{publicstaticvoidmain(String[]args)throwsException{SimpleDateFormatsdfnewSimpleDateFormat(HH:mm:ss.SSS);// 1. Push 消费者消费组 lesson06_consumer_groupDefaultMQPushConsumerconsumernewDefaultMQPushConsumer(lesson06_consumer_group);consumer.setNamesrvAddr(127.0.0.1:9876);// 2. 订阅 TopicLesson06 全部消息* tag 不过滤consumer.subscribe(TopicLesson06,*);// 3. 注册监听器每收到一条打印接收时间 内容 Broker 落盘时间AtomicIntegerreceivedCountnewAtomicInteger(0);consumer.registerMessageListener(newMessageListenerConcurrently(){OverridepublicConsumeConcurrentlyStatusconsumeMessage(ListMessageExtmsgs,ConsumeConcurrentlyContextcontext){for(MessageExtmsg:msgs){StringbodynewString(msg.getBody(),StandardCharsets.UTF_8);System.out.println(sdf.format(newDate()) 收到消息: bodytagmsg.getTags()Broker落盘sdf.format(newDate(msg.getStoreTimestamp())));receivedCount.incrementAndGet();}returnConsumeConcurrentlyStatus.CONSUME_SUCCESS;}});consumer.start();System.out.println(消费者已启动正在等待消息……预计 2 条消息A 立刻到、消息B 约 10~20 秒后到收满 2 条或等满 50 秒自动退出);// 4. 收满 2 条或最多等 50 秒消息B 要走回查留足时间longdeadlineSystem.currentTimeMillis()50_000;while(receivedCount.get()2System.currentTimeMillis()deadline){Thread.sleep(500);}consumer.shutdown();System.out.println(sdf.format(newDate()) 观察结束消费者已关闭共收到 receivedCount.get() 条。);}}消费者代码本身和普通消息消费几乎一样——这正是事务消息的设计目标消费者无感知。它根本不知道哪条消息走过半消息、哪条走过回查只看到最终 COMMIT 的消息。五、实测输出两条消息的时间线对比以下是本教程同款环境Windows 11 Docker 29 JDK 17下的真实运行输出时间戳按程序实际打印整理队列编号每次运行可能不同属正常。生产者输出18:24:40.612 事务生产者已启动。 18:24:40.830 [本地事务] 消息A 执行成功模拟订单已落库返回 COMMIT 18:24:40.833 发送消息A 完成SendResult [sendStatusSEND_OK, msgId7F000001..., messageQueueMessageQueue [topicTopicLesson06, brokerNamebroker-a, queueId2], queueOffset3] 18:24:40.837 [本地事务] 消息B 结果未知模拟本地操作超时返回 UNKNOWN等待 Broker 回查 18:24:40.837 发送消息B 完成SendResult [sendStatusSEND_OK, msgId7F000001..., messageQueueMessageQueue [topicTopicLesson06, brokerNamebroker-a, queueId3], queueOffset4] 生产者发送完毕保持 JVM 存活 30 秒等待 Broker 回查消息B…… 已等待 5 秒…… 18:24:47.712 [回查] 收到 Broker 回查消息消息B用户下单成功请给用户发短信通知 → 查库后判定成功返回 COMMIT 已等待 10 秒…… 已等待 15 秒…… 已等待 20 秒…… 已等待 25 秒…… 已等待 30 秒…… 18:25:10.876 生产者已关闭。消费者输出消费者已启动正在等待消息……预计 2 条消息A 立刻到、消息B 约 10~20 秒后到收满 2 条或等满 50 秒自动退出 18:24:40.852 收到消息: 消息A用户下单成功请给用户加 100 积分tagTagCommitBroker落盘18:24:40.837 18:24:47.724 收到消息: 消息B用户下单成功请给用户发短信通知tagTagCheckbackBroker落盘18:24:47.714 18:24:48.008 观察结束消费者已关闭共收到 2 条。看懂这组输出这就是本课的全部精华时间点事件说明18:24:40.83消息 A 本地事务 COMMIT发送后 15 毫秒消费者就收到了——正常提交路径秒级可见18:24:40.84消息 B 本地事务返回 UNKNOWN此时消息 B 是半消息消费者还看不到它18:24:47.71收到 Broker 回查距发送约 7 秒本教程把回查间隔配成了 10 秒实际在 10 秒窗口内触发Broker 上门问消息 B 成了没18:24:47.71消息 B 才正式落盘看消费者那行的Broker落盘18:24:47.714——消息 B 直到回查这一刻才真正生效比消息 A 晚了 7 秒一句话总结消息 A 走发完即提交消息 B 走半消息挂起 → 回查 → 提交两条消息最终都被消费者收到本地事务和发消息实现了最终一致。六、事务消息 vs 2PC vs 本地事务表怎么选和 2PC 的区别最终一致 vs 强一致听到分布式事务你可能听过 2PC两阶段提交。两者完全是两种思路事务消息本课2PC一致性最终一致秒级延迟后两边对齐强一致事务结束时所有方同时知道结果锁资源不锁任何业务资源全程锁定资源等大家投票是否阻塞不阻塞主流程用户立刻走人协调者挂了所有参与者全卡住类比网购下单后你该干嘛干嘛卖家慢慢发货当面交易一手交钱一手交货谁都别走适合场景下单送积分、下单发短信通知跨行转账这种必须同时成功/同时失败的核心账务一句话事务消息是弱一致但好用的方案绝大多数互联网业务场景用它就够了。和本地事务表方案的区别还有一种常见做法不依赖 MQ 的事务消息功能写业务库时在同一个数据库事务里顺带往一张消息表插一条待发送记录再起一个定时任务每隔几秒扫这张表把待发送的记录发到 MQ发成功就标记掉。RocketMQ 事务消息本地事务表 定时扫表业务侵入小写一个事务监听器就行大建表、写扫表任务、写补偿逻辑实时性半秒级立刻可见取决于扫表间隔常见分钟级一致性保证MQ 引擎 回查机制兜底全靠你自己的代码写对换 MQ换别的 MQ如 Kafka这套就没了任何 MQ 都能用选型一句话用 RocketMQ 就优先事务消息要是你用的 MQ 不支持事务消息再退而求其次上本地事务表。七、三个最容易踩的坑坑 1生产者发完就退出消息 B 永远收不到。回查checkLocalTransaction是 Broker主动回调生产者进程的。生产者一退出10 秒后 Broker 上门找不到人消息 B 就永远卡在半消息状态。Demo 里发完消息后保持 JVM 存活 30 秒就是这个原因。真实项目里生产者是常驻服务没这个问题。坑 2以为 UNKNOWN 就是消息丢了。本地事务返回 UNKNOWN 后立刻去看 Topic发现消费者没有消息以为发失败了。UNKNOWN 的含义是结果待定不是失败。Broker 会挂起这条半消息等回查。要耐心等回查日志打出来才是转折点。坑 3先起生产者、再起消费者一条都收不到。新消费组默认从最新位置开始消费启动之前的历史消息不算。而且新 Topic 要等客户端路由刷新默认 30 秒才能被发现。所以一定要先起消费者、再发生产者。八、挑战题答案在仓库跑起来才知道⭐ 把checkLocalTransaction的返回值从COMMIT_MESSAGE改成ROLLBACK_MESSAGE重新运行观察消费者是否只收到消息 A、消息 B 被丢弃回查判定失败。⭐⭐ 把executeLocalTransaction里消息 B 的返回值也改成ROLLBACK_MESSAGE本地事务直接失败观察消费者是否只收到 1 条消息。想一想这时还会触发回查吗为什么⭐⭐⭐ 升级成真回查——用一个ConcurrentHashMap模拟业务库executeLocalTransaction对消息 B 故意不写结果模拟宕机checkLocalTransaction里才去 Map 查并补写结果再 COMMIT。还原本地事务超时后靠回查救场的真实场景。九、面试回答模板面试官事务消息解决什么问题本地数据库事务和发 MQ 消息这两件事跨系统、没法用同一个事务包起来事务消息用半消息 本地事务 回查实现最终一致。见本文第二节、第三节追问Half 半消息是什么为什么消费者一开始看不到已写到 Broker、但存在内部 Half Topic 里的消息对消费者完全不可见相当于冻结在担保交易里的定金。见本文第二节追问本地事务返回 UNKNOWN 后会发生什么回查机制怎么兜底Broker 每隔transactionCheckInterval回调生产者checkLocalTransaction真实项目里就是查数据库判断成没成默认回查 15 次。见本文第三节追问事务消息和 2PC 有什么区别事务消息是最终一致、不锁资源、不阻塞2PC 是强一致、全程锁资源、协调者挂了会阻塞所有方。绝大多数互联网业务场景用事务消息就够了。见本文第六节关于这个系列本文是「Java 后端实战精通营」系列第 6 篇原则实战驱动、由浅到深、面试向每篇文章的结论都可以亲手验证。RocketMQ 实战精通营10 课https://gitee.com/j67mk2/rocketmq-journey本文对应源码位置lesson-06/事务消息生产者 消费者完整工程演示 COMMIT 与回查两条路径系列文章一览按发布顺序篇主题1RocketMQ 入门Docker 一行起三件套跑通你的第一条消息2RocketMQ 发送方式同步异步批量单向消息都怎么发出去3RocketMQ 消费模式集群、广播与重试消息怎么被吃掉4RocketMQ 可靠性发送重试加幂等消息一条都不丢5RocketMQ 顺序消息订单流程不乱套的秘密6RocketMQ 事务消息订单与积分的最终一致7RocketMQ 延迟消息30 分钟未支付自动关单怎么做8RocketMQ 积压治理百万消息堵在队列怎么办9RocketMQ 集群高可用与过滤主从架构 Tag 精准投递10RocketMQ 面试冲刺高频考点一口气背完下一篇预告《RocketMQ 延迟消息30 分钟未支付自动关单怎么做》——下单后 30 分钟没付款要自动关单怎么实现RocketMQ 的延迟消息是怎么做到的跑完有任何报错把终端输出发评论区一起排查。
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →