SpringCloud集成RocketMQ实现事务消息方案
前边的话当前SpringCloud作为微服务开发的首选开源方案提供了完善的微服务开发技术套件不过针对分布式领域的难题–分布式事务控制并没有成熟的方案本篇将介绍作为柔性事务控制的优秀方案RocketMQ的使用原理和方法。通过本案例的学习掌握SpringCloud集成RocketMQ事务分布式事务控制的方法。分布式事务控制系列文章分布式事务视频教程下载RocketMQ事务消息方案RocketMQ 是一个来自阿里巴巴的分布式消息中间件于 2012 年开源并在 2017 年正式成为 Apache 顶级项目。据了解包括阿里云上的消息产品以及收购的子公司在内阿里集团的消息产品全线都运行在 RocketMQ 之上并且最近几年的双十一大促中RocketMQ 都有抢眼表现。Apache RocketMQ 4.3之后的版本正式支持事务消息为分布式事务实现提供了便利性支持。RocketMQ 事务消息设计则主要是为了解决 Producer 端的消息发送与本地事务执行的原子性问题RocketMQ 的设计中 broker 与 producer 端的双向通信能力使得 broker 天生可以作为一个事务协调者存在而 RocketMQ 本身提供的存储机制为事务消息提供了持久化能力RocketMQ 的高可用机制以及可靠消息设计则为事务消息在系统发生异常时依然能够保证达成事务的最终一致性。在RocketMQ 4.3后实现了完整的事务消息实际上其实是对本地消息表的一个封装将本地消息表移动到了MQ内部解决 Producer 端的消息发送与本地事务执行的原子性问题。执行流程如下为方便理解我们还以注册送积分的例子来描述 整个流程。Producer 即MQ发送方本例中是用户服务负责新增用户。MQ订阅方即消息消费方本例中是积分服务负责新增积分。1、Producer 发送事务消息Producer MQ发送方发送事务消息至MQ ServerMQ Server将消息状态标记为Prepared预备状态注意此时这条消息消费者MQ订阅方是无法消费到的。本例中Producer 发送 ”增加积分消息“ 到MQ Server。2、MQ Server回应消息发送成功MQ Server接收到Producer 发送给的消息则回应发送成功表示MQ已接收到消息。3、Producer 执行本地事务Producer 端执行业务代码逻辑通过本地数据库事务控制。本例中Producer 执行添加用户操作。4、消息投递若Producer 本地事务执行成功则自动向MQServer发送commit消息MQ Server接收到commit消息后将”增加积分消息“ 状态标记为可消费此时MQ订阅方积分服务即正常消费消息若Producer 本地事务执行失败则自动向MQServer发送rollback消息MQ Server接收到rollback消息后 将删除”增加积分消息“ 。MQ订阅方积分服务消费消息消费成功则向MQ回应ack否则将重复接收消息。这里ack默认自动回应即程序执行正常则自动回应ack。5、事务回查如果执行Producer端本地事务过程中执行端挂掉或者超时MQ Server将会不停的询问同组的其他 Producer来获取事务执行状态这个过程叫事务回查。MQ Server会根据事务回查结果来决定是否投递消息。以上主干流程已由RocketMQ实现对用户侧来说用户需要分别实现本地事务执行以及本地事务回查方法因此只需关注本地事务的执行状态即可。RoacketMQ提供RocketMQLocalTransactionListener接口public interface RocketMQLocalTransactionListener { /** - 发送prepare消息成功此方法被回调该方法用于执行本地事务 - param msg 回传的消息利用transactionId即可获取到该消息的唯一Id - param arg 调用send方法时传递的参数当send时候若有额外的参数可以传递到send方法中这里能获取到 - return 返回事务状态COMMIT提交 ROLLBACK回滚 UNKNOW回调 */ RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg); /** - param msg 通过获取transactionId来判断这条消息的本地事务执行状态 - return 返回事务状态COMMIT提交 ROLLBACK回滚 UNKNOW回调 */ RocketMQLocalTransactionState checkLocalTransaction(Message msg); }发送事务消息以下是RocketMQ提供用于发送事务消息的APITransactionMQProducer producer new TransactionMQProducer(ProducerGroup); producer.setNamesrvAddr(127.0.0.1:9876); producer.start(); //设置TransactionListener实现 producer.setTransactionListener(transactionListener //发送事务消息 SendResult sendResult producer.sendMessageInTransaction(msg, null);案例说明本实例通过RocketMQ中间件实现可靠消息最终一致性分布式事务模拟两个账户的转账交易过程。两个账户在分别在不同的银行(张三在bank1、李四在bank2)bank1、bank2是两个微服务。交易过程是张三给李四转账指定金额。上述交易步骤张三扣减金额与给bank2发转账消息两个操作必须是一个整体性的事务。案例组成数据库MySQL-5.7.25包括bank1和bank2两个数据库。JDK64位 jdk1.8.0_201rocketmq 服务端RocketMQ-4.5.0rocketmq 客户端RocketMQ-Spring-Boot-starter.2.0.2-RELEASE微服务框架spring-boot-2.1.3、spring-cloud-Greenwich.RELEASE微服务及数据库的关系 dtx/dtx-txmsg-demo/dtx-txmsg-demo-bank1 银行1操作张三账户 连接数据库bank1 dtx/dtx-txmsg-demo/dtx-txmsg-demo-bank2 银行2操作李四账户连接数据库bank2本示例程序技术架构如下交互流程如下1、Bank1向MQ Server发送转账消息2、Bank1执行本地事务扣减金额3、Bank2接收消息执行本地事务添加金额创建数据库导入数据库脚本sql\bank1.sql、sql\bank2.sql已经导过不用重复导入。脚本https://github.com/pbteach/pbdtx/tree/master/sql创建bank1库并导入以下表结构和数据(包含张三账户)CREATE DATABASE bank1 CHARACTER SET utf8 COLLATE utf8_general_ci; DROP TABLE IF EXISTS account_info; CREATE TABLE account_info ( id bigint(20) NOT NULL AUTO_INCREMENT, account_name varchar(100) CHARACTER SET utf8 COLLATE utf8_bin NULL DEFAULT NULL COMMENT 户主姓名, account_no varchar(100) CHARACTER SET utf8 COLLATE utf8_bin NULL DEFAULT NULL COMMENT 银行卡号, account_password varchar(100) CHARACTER SET utf8 COLLATE utf8_bin NULL DEFAULT NULL COMMENT 帐户密码, account_balance double NULL DEFAULT NULL COMMENT 帐户余额, PRIMARY KEY (id) USING BTREE ) ENGINE InnoDB AUTO_INCREMENT 5 CHARACTER SET utf8 COLLATE utf8_bin ROW_FORMAT Dynamic; INSERT INTO account_info VALUES (2, 张三的账户, 1, , 10000);创建bank2库并导入以下表结构和数据(包含李四账户)CREATE DATABASE bank2 CHARACTER SET utf8 COLLATE utf8_general_ci; CREATE TABLE account_info ( id bigint(20) NOT NULL AUTO_INCREMENT, account_name varchar(100) CHARACTER SET utf8 COLLATE utf8_bin NULL DEFAULT NULL COMMENT 户主姓名, account_no varchar(100) CHARACTER SET utf8 COLLATE utf8_bin NULL DEFAULT NULL COMMENT 银行卡号, account_password varchar(100) CHARACTER SET utf8 COLLATE utf8_bin NULL DEFAULT NULL COMMENT 帐户密码, account_balance double NULL DEFAULT NULL COMMENT 帐户余额, PRIMARY KEY (id) USING BTREE ) ENGINE InnoDB AUTO_INCREMENT 5 CHARACTER SET utf8 COLLATE utf8_bin ROW_FORMAT Dynamic; INSERT INTO account_info VALUES (3, 李四的账户, 2, NULL, 0);在bank1、bank2数据库中新增de_duplication交易记录表(去重表)用于交易幂等控制。DROP TABLE IF EXISTS de_duplication; CREATE TABLE de_duplication ( tx_no varchar(64) COLLATE utf8_bin NOT NULL, create_time datetime(0) NULL DEFAULT NULL, PRIMARY KEY (tx_no) USING BTREE ) ENGINE InnoDB CHARACTER SET utf8 COLLATE utf8_bin ROW_FORMAT Dynamic;启动RocketMQ1下载RocketMQ服务器下载地址http://mirrors.tuna.tsinghua.edu.cn/apache/rocketmq/4.5.0/rocketmq-all-4.5.0-bin-release.zip2解压并启动启动nameserver:set ROCKETMQ_HOME[rocketmq服务端解压路径] start [rocketmq服务端解压路径]/bin/mqnamesrv.cmd启动broker:set ROCKETMQ_HOME[rocketmq服务端解压路径] start [rocketmq服务端解压路径]/bin/mqbroker.cmd -n 127.0.0.1:9876 autoCreateTopicEnabletrue案例工程案例代码https://github.com/pbteach/pbdtx/tree/master/dtx-txmsg-demo两个测试工程如下dtx/dtx-txmsg-demo/dtx-txmsg-demo-bank1 操作张三账户连接数据库bank1 dtx/dtx-txmsg-demo/dtx-txmsg-demo-bank2 操作李四账户连接数据库bank21父工程maven依赖说明在dtx父工程中指定了SpringBoot和SpringCloud版本groupIdcom.pbteach.dtx/groupId artifactIddtx-parent/artifactId packagingpom/packaging version1.0-SNAPSHOT/version dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-dependencies/artifactId version2.1.3.RELEASE/version typepom/type scopeimport/scope /dependency dependency groupIdorg.springframework.cloud/groupId artifactIdspring-cloud-dependencies/artifactId versionGreenwich.RELEASE/version typepom/type scopeimport/scope /dependency在dtx-txmsg-demo父工程中指定了rocketmq-spring-boot-starter的版本。dependency groupIdorg.apache.rocketmq/groupId artifactIdrocketmq-spring-boot-starter/artifactId version2.0.2/version /dependency2配置rocketMQ在application-local.propertis中配置rocketMQ nameServer地址及生产组rocketmq.producer.group producer_bank2 rocketmq.name-server 127.0.0.1:9876dtx-txmsg-demo-bank1代码https://github.com/pbteach/pbdtx/tree/master/dtx-txmsg-demo/dtx-txmsg-demo-bank1dtx-txmsg-demo-bank1实现如下功能1、张三扣减金额提交本地事务。2、向MQ发送转账消息。2DaoMapper Component public interface AccountInfoDao { Update(update account_info set account_balanceaccount_balance#{amount} where account_no#{accountNo}) int updateAccountBalance(Param(accountNo) String accountNo, Param(amount) Double amount); Select(select count(1) from de_duplication where tx_no #{txNo}) int isExistTx(String txNo); Insert(insert into de_duplication values(#{txNo},now());) int addTx(String txNo); }3AccountInfoServiceService Slf4j public class AccountInfoServiceImpl implements AccountInfoService { Resource private RocketMQTemplate rocketMQTemplate; Autowired private AccountInfoDao accountInfoDao; /** * 更新帐号余额-发送消息 * producer向MQ Server发送消息 * * param accountChangeEvent */ Override public void sendUpdateAccountBalance(AccountChangeEvent accountChangeEvent) { //构建消息体 JSONObject jsonObject new JSONObject(); jsonObject.put(accountChange,accountChangeEvent); MessageString message MessageBuilder.withPayload(jsonObject.toJSONString()).build(); TransactionSendResult sendResult rocketMQTemplate.sendMessageInTransaction(producer_group_txmsg_bank1, topic_txmsg, message, null); log.info(send transcation message body{},result{},message.getPayload(),sendResult.getSendStatus()); } /** * 更新帐号余额-本地事务 * producer发送消息完成后接收到MQ Server的回应即开始执行本地事务 * * param accountChangeEvent */ Transactional Override public void doUpdateAccountBalance(AccountChangeEvent accountChangeEvent) { log.info(开始更新本地事务事务号{},accountChangeEvent.getTxNo()); accountInfoDao.updateAccountBalance(accountChangeEvent.getAccountNo(),accountChangeEvent.getAmount() * -1); //为幂等作准备 accountInfoDao.addTx(accountChangeEvent.getTxNo()); if(accountChangeEvent.getAmount() 2){ throw new RuntimeException(bank1更新本地事务时抛出异常); } log.info(结束更新本地事务事务号{},accountChangeEvent.getTxNo()); } }4RocketMQLocalTransactionListener编写RocketMQLocalTransactionListener接口实现类实现执行本地事务和事务回查两个方法。Component Slf4j RocketMQTransactionListener(txProducerGroup producer_group_txmsg_bank1) public class ProducerTxmsgListener implements RocketMQLocalTransactionListener { Autowired AccountInfoService accountInfoService; Autowired AccountInfoDao accountInfoDao; //消息发送成功回调此方法此方法执行本地事务 Override Transactional public RocketMQLocalTransactionState executeLocalTransaction(Message message, Object arg) { //解析消息内容 try { String jsonString new String((byte[]) message.getPayload()); JSONObject jsonObject JSONObject.parseObject(jsonString); AccountChangeEvent accountChangeEvent JSONObject.parseObject(jsonObject.getString(accountChange), AccountChangeEvent.class); //扣除金额 accountInfoService.doUpdateAccountBalance(accountChangeEvent); return RocketMQLocalTransactionState.COMMIT; } catch (Exception e) { log.error(executeLocalTransaction 事务执行失败,e); e.printStackTrace(); return RocketMQLocalTransactionState.ROLLBACK; } } //此方法检查事务执行状态 Override public RocketMQLocalTransactionState checkLocalTransaction(Message message) { RocketMQLocalTransactionState state; final JSONObject jsonObject JSON.parseObject(new String((byte[]) message.getPayload())); AccountChangeEvent accountChangeEvent JSONObject.parseObject(jsonObject.getString(accountChange),AccountChangeEvent.class); //事务id String txNo accountChangeEvent.getTxNo(); int isexistTx accountInfoDao.isExistTx(txNo); log.info(回查事务事务号: {} 结果: {}, accountChangeEvent.getTxNo(),isexistTx); if(isexistTx0){ state RocketMQLocalTransactionState.COMMIT; }else{ state RocketMQLocalTransactionState.UNKNOWN; } return state; } }5ControllerRestController Slf4j public class AccountInfoController { Autowired private AccountInfoService accountInfoService; GetMapping(value /transfer) public String transfer(RequestParam(accountNo)String accountNo,RequestParam(amount) Double amount){ String tx_no UUID.randomUUID().toString(); AccountChangeEvent accountChangeEvent new AccountChangeEvent(accountNo,amount,tx_no); accountInfoService.sendUpdateAccountBalance(accountChangeEvent); return 转账成功; } }dtx-txmsg-demo-bank2代码https://github.com/pbteach/pbdtx/tree/master/dtx-txmsg-demo/dtx-txmsg-demo-bank2dtx-txmsg-demo-bank2需要实现如下功能1、监听MQ接收消息。2、接收到消息增加账户金额。1 Service注意为避免消息重复发送这里需要实现幂等。Service Slf4j public class AccountInfoServiceImpl implements AccountInfoService { Autowired AccountInfoDao accountInfoDao; /** * 消费消息更新本地事务添加金额 * param accountChangeEvent */ Override Transactional public void addAccountInfoBalance(AccountChangeEvent accountChangeEvent) { log.info(bank2更新本地账号账号{},金额{},accountChangeEvent.getAccountNo(),accountChangeEvent.getAmount()); //幂等校验 int existTx accountInfoDao.isExistTx(accountChangeEvent.getTxNo()); if(existTx0){ //执行更新 accountInfoDao.updateAccountBalance(accountChangeEvent.getAccountNo(),accountChangeEvent.getAmount()); //添加事务记录 accountInfoDao.addTx(accountChangeEvent.getTxNo()); log.info(更新本地事务执行成功本次事务号: {}, accountChangeEvent.getTxNo()); }else{ log.info(更新本地事务执行失败本次事务号: {}, accountChangeEvent.getTxNo()); } } }2MQ监听类Component RocketMQMessageListener(topic topic_txmsg,consumerGroup consumer_txmsg_group_bank2) Slf4j public class TxmsgConsumer implements RocketMQListenerString { Autowired AccountInfoService accountInfoService; Override public void onMessage(String s) { log.info(开始消费消息:{},s); //解析消息为对象 final JSONObject jsonObject JSON.parseObject(s); AccountChangeEvent accountChangeEvent JSONObject.parseObject(jsonObject.getString(accountChange),AccountChangeEvent.class); //调用service增加账号金额 accountChangeEvent.setAccountNo(2); accountInfoService.addAccountInfoBalance(accountChangeEvent); } }测试场景bank1本地事务失败则bank1不发送转账消息。bank2接收转账消息失败会进行重试发送消息。bank2多次消费同一个消息实现幂等。视频分享https://blog.csdn.net/weixin_44062339/article/details/101010064

相关新闻

Android Activity结构 Activity view window 以及xml布局文件之间的关系

Android Activity结构 Activity view window 以及xml布局文件之间的关系

Activity中onCreate方法中,通过setContentView方法把布局文件加载进去 Override protected void onCreate(Bundle savedInstanceState) {super.onCreate(savedInstanceState);//加载布局文件setContentView(R.layout.activity_path_measure); }

2026/7/28 18:14:08 阅读更多 →
一份“热腾腾”的面经分享(涵盖国内知名大厂)

一份“热腾腾”的面经分享(涵盖国内知名大厂)

百度现场面试:JVM算法Redis数据库!(三面) 百度一面(现场) 自我介绍Java中的多态为什么要同时重写hashcode和equalsHashmap的原理Hashmap如何变线程安全,每种方式的优缺点垃圾回收机制Jvm的参数你知道的说…

2026/7/28 18:14:08 阅读更多 →
ComfyUI-VideoHelperSuite VHS_VideoCombine节点缺失问题:完整故障排查与解决方案

ComfyUI-VideoHelperSuite VHS_VideoCombine节点缺失问题:完整故障排查与解决方案

ComfyUI-VideoHelperSuite VHS_VideoCombine节点缺失问题:完整故障排查与解决方案 【免费下载链接】ComfyUI-VideoHelperSuite Nodes related to video workflows 项目地址: https://gitcode.com/gh_mirrors/co/ComfyUI-VideoHelperSuite ComfyUI-VideoHelpe…

2026/7/28 18:14:08 阅读更多 →

最新新闻

如何用AI图层分离技术5分钟完成专业设计工作

如何用AI图层分离技术5分钟完成专业设计工作

如何用AI图层分离技术5分钟完成专业设计工作 【免费下载链接】layerdivider A tool to divide a single illustration into a layered structure. 项目地址: https://gitcode.com/gh_mirrors/la/layerdivider 在数字设计领域,图层分离一直是一项耗时且技术性…

2026/7/28 18:23:11 阅读更多 →
国家中小学智慧教育平台电子课本下载完整指南:三步获取所有教材PDF

国家中小学智慧教育平台电子课本下载完整指南:三步获取所有教材PDF

国家中小学智慧教育平台电子课本下载完整指南:三步获取所有教材PDF 【免费下载链接】tchMaterial-parser 国家中小学智慧教育平台 电子课本下载工具,帮助您从智慧教育平台中获取电子课本的 PDF 文件网址并进行下载,让您更方便地获取课本内容。…

2026/7/28 18:23:11 阅读更多 →
300. 最长上升子序列

300. 最长上升子序列

思路一:动态规划 建立dp表,dp[i]表示含第i个数字的最长上升子序列的长度 求dp[i]时,向前遍历找出比i元素小的元素j,则动态方程为dp[i] max(dp[i],dp[j] 1) class Solution(object):def lengthOfLIS(self, nums):size len(nums)…

2026/7/28 18:23:11 阅读更多 →
5分钟找回QQ空间全部历史说说的终极指南:GetQzonehistory使用教程

5分钟找回QQ空间全部历史说说的终极指南:GetQzonehistory使用教程

5分钟找回QQ空间全部历史说说的终极指南:GetQzonehistory使用教程 【免费下载链接】GetQzonehistory 获取QQ空间发布的历史说说 项目地址: https://gitcode.com/GitHub_Trending/ge/GetQzonehistory 你是否曾想找回那些被时间淹没的QQ空间历史说说&#xff1…

2026/7/28 18:23:11 阅读更多 →
JSP总结

JSP总结

JSP总结 文章目录JSP总结一.jsp工作原理和生命周期二.jsp内置对象三.incude指令和include行为四.jsp作用域五.cookie和session一.jsp工作原理和生命周期 工作原理: 执行过程以hello.jsp为例: 把 hello.jsp转译为hello_jsp.javahello_jsp.java 位于 d:…

2026/7/28 18:23:10 阅读更多 →
EDA 工程师竞业锁定、空窗 2 年:合规赚钱完整方案(零违约风险,杜绝赔付违约金)

EDA 工程师竞业锁定、空窗 2 年:合规赚钱完整方案(零违约风险,杜绝赔付违约金)

EDA 工程师竞业锁定、空窗 2 年:合规赚钱完整方案(零违约风险,杜绝赔付违约金)先确立核心法理底线: 竞业限制禁止三类行为(无论全职 / 兼职 / 外包 / 顾问 / 项目合作)入职存在实质竞争关系的 E…

2026/7/28 18:22:10 阅读更多 →

日新闻

告别臃肿!3步让你的暗影精灵笔记本重获新生

告别臃肿!3步让你的暗影精灵笔记本重获新生

告别臃肿!3步让你的暗影精灵笔记本重获新生 【免费下载链接】OmenSuperHub Control Omen laptop performance, fan speeds, and keyboard lighting, and unlock power limits. 项目地址: https://gitcode.com/gh_mirrors/om/OmenSuperHub 你是否也曾为官方Om…

2026/7/28 0:00:43 阅读更多 →
RAG必踩坑!财报法规检索不准?这款开源工具让答案浮出水面,准确率飙升98.7%!

RAG必踩坑!财报法规检索不准?这款开源工具让答案浮出水面,准确率飙升98.7%!

做 RAG 的人应该都踩过这个致命的坑:把几百页的财报、法规、技术手册扔给向量库,问一个具体问题,搜出来的全是沾边但没用的内容 —— 关键信息要么被硬切块拆碎了,要么藏在几十条结果的最下面。语义相似≠真正相关,这个…

2026/7/28 0:00:43 阅读更多 →
抖音视频文案提取工具全指南:免费2026版、手机App、在线工具一网打尽

抖音视频文案提取工具全指南:免费2026版、手机App、在线工具一网打尽

2026年做短视频运营,从抖音上扒文案早就不是偷偷抄笔记的事了。我刚开始做内容的时候,每天刷半小时抖音,手动把爆款视频的口播敲进备忘录,一条2分钟的视频得花十来分钟,碰到语速快的还要反复回听。后来试了一圈工具&am…

2026/7/28 0:00:43 阅读更多 →

周新闻

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 数据集6000张 完整源码已标注数据集训练好的模型环境配置教程程序运行说明文档,可以直接使用!系统支持图片、视频、摄像头等多种方式检测裂缝,功能强大实用。 1数据集6000张 8各类别

2026/7/28 12:04:22 阅读更多 →
深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

pubg数据集 精选原图1.42万数据 1.49万标签 无任何重复、算法增强或冗余图像! pubg绝地求生目标检测数据集 1分类:e_body,14905个标签,txt格式 共计14244张图,99%为640*640尺寸图像 适合yolo目标检测、AI训练关键词&am…

2026/7/28 8:29:16 阅读更多 →
Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex检测数据集数据集详情检测类别: allies enemy tag图片总量:7247张训练集:5139张验证集:1425张测试集:683张标注状态:全部已标注,即拿即用数据格式:支持YOLO格式及其他格式&#…

2026/7/28 5:03:42 阅读更多 →

月新闻