Java高级全套教程(十二)—— 分布式事务超详细实战全解(基础理论+双方案企业级实战)
第一章 数据库本地事务核心原理
1.1 事务核心定义
事务是数据库层面保障数据操作可靠性的核心机制,指一组关联性的SQL操作集合,集合内的所有数据库操作是一个不可分割的原子单元。事务执行遵循“全员成功、全员回滚”原则:若所有SQL执行无异常,整体提交生效;若任意一条SQL执行失败、程序报错、网络中断,所有已执行的操作全部撤销,数据恢复至操作前状态。
事务的核心作用是解决业务操作中数据一致性问题,避免多步数据库操作出现“部分成功、部分失败”的数据错乱场景。
业务场景举例:用户下单扣库存、生成订单、扣余额是一组关联业务,任意一步失败,所有操作都必须回滚,不能出现“库存扣除但订单创建失败”的脏数据。
1.2 本地事务概念
本地事务又称数据库事务,是依托关系型数据库(MySQL、Oracle等)原生事务机制实现的事务控制。其核心特征为:事务涉及的所有业务数据、数据库操作,均在同一个数据库、同一个服务节点内完成,无需跨服务、跨库、跨节点协作。
本地事务仅适用于单体架构系统,依托数据库ACID特性即可完美保障数据一致性;在微服务分布式架构中,多服务、多数据库拆分后,本地事务将完全失效,无法解决跨服务数据一致性问题。
1.3 事务四大核心特性(ACID)
关系型数据库事务具备四大刚性特性,是所有数据一致性保障的基础:
-
原子性(Atomicity):事务是最小执行单元,不可拆分。所有操作要么全部执行成功提交,要么全部失败回滚,不存在中间状态。
-
一致性(Consistency):事务执行前后,数据库整体数据完整性、业务约束始终保持一致。例如转账场景,总资金总额不会发生变化。
-
隔离性(Isolation):多个并发执行的事务相互隔离、互不干扰,数据库通过不同隔离级别控制并发事务的干扰程度,避免脏数据产生。
-
持久性(Durability):事务一旦提交成功,数据变更会永久写入磁盘,即使服务器宕机、重启,数据也不会丢失。
第二章 并发事务引发的核心问题
数据库支持多事务并发执行,高并发场景下,多个事务同时操作同一批数据,若没有合理的隔离机制,会引发数据异常问题。核心分为三类读写冲突问题、一类写冲突问题,是定义事务隔离级别的核心依据。
2.1 脏写(写-写冲突)
定义:两个及以上事务同时操作同一行数据,后提交的事务覆盖先提交事务的修改结果,导致前序事务的数据修改丢失。
场景复现:事务A修改用户余额未提交,事务B同时修改同一用户余额并提交,事务A后续提交后,直接覆盖B的修改数据,造成数据丢失。
解决方案:强制事务串行执行,同一数据同一时间仅允许一个事务执行写操作,杜绝并行写冲突。
2.2 脏读(读-写冲突)
定义:一个事务读取到了另一个事务未提交的临时数据,后续若该事务回滚,当前读取的数据即为无效脏数据。
场景复现:事务A扣除商品库存(未提交),事务B读取到库存减少的数据并执行后续业务,事务A因异常回滚恢复库存,导致事务B基于脏数据执行业务,引发数据错乱。
解决方案:实现写优先机制,数据写操作未提交完成前,禁止其他事务读取数据。
2.3 不可重复读(读-写冲突)
定义:同一事务内,两次读取同一数据结果不一致。事务第一次读取数据后,其他事务修改并提交该数据,导致当前事务再次查询时数据发生变化。
核心区别:脏读是读取未提交数据,不可重复读是读取已提交的修改数据,聚焦数据更新场景。
解决方案:实现读优先机制,事务读取数据期间,禁止其他事务修改、更新数据。
2.4 幻读(读-写冲突)
定义:同一事务内,根据相同查询条件多次查询,查询结果的数据集数量发生变化。其他事务新增/删除符合条件的数据,导致当前事务出现“幻觉数据”。
核心区别:不可重复读针对单条数据修改,幻读针对数据集新增/删除。
解决方案:锁定查询范围,事务读取期间,禁止其他事务新增、删除对应条件的数据。
第三章 MySQL事务隔离级别详解
MySQL InnoDB引擎实现了SQL标准的4种事务隔离级别,隔离级别从低到高,并发性能逐步降低,数据一致性逐步提升,可根据业务场景灵活选型。
3.1 四大隔离级别特性对照表
| 读未提交(Read Uncommitted) | 存在 | 存在 | 存在 | 极少使用,仅用于数据实时监控场景 |
| 读已提交(Read Committed) | 杜绝 | 存在 | 存在 | Oracle、SQL Server默认级别,适用于大多数并发业务 |
| 可重复读(Repeatable Read) | 杜绝 | 杜绝 | 存在(理论) | MySQL默认级别,适配绝大多数企业业务场景 |
| 串行化(Serializable) | 杜绝 | 杜绝 | 杜绝 | 超高数据一致性场景,无并发需求,性能极低 |
3.2 各级别详细说明
3.2.1 读未提交
最低隔离级别,允许事务读取其他事务未提交的修改数据,无法杜绝任何并发问题,数据一致性极差,生产环境基本不使用。
3.2.2 读已提交
仅允许读取其他事务已提交的数据,彻底解决脏读问题。但同一事务内多次读取同一数据,可能读取到其他事务已提交的更新数据,存在不可重复读、幻读问题。
3.2.3 可重复读
MySQL默认事务隔离级别,保证同一事务多次读取同一数据结果完全一致,彻底解决脏读、不可重复读问题。InnoDB引擎通过MVCC机制极大缓解幻读问题,满足99%的企业业务需求。
3.2.4 串行化
最高隔离级别,强制所有事务串行顺序执行,完全规避所有并发问题。但会产生大量锁等待、锁竞争,并发性能极低,仅用于金融、支付等对数据一致性要求极致严苛的场景。
第四章 分布式事务核心认知与解决方案选型
4.1 分布式事务产生背景
单体架构中,所有业务、数据集中在一个数据库,本地事务可完美保障一致性。微服务架构下,业务被拆分为多个独立微服务,每个服务对应独立数据库,跨服务业务操作(下单、扣库存、扣款)会涉及多库、多服务操作,本地事务失效,由此产生分布式事务问题。
分布式事务核心痛点:跨服务操作无法实现原子性,容易出现部分服务执行成功、部分失败,导致数据不一致。
4.2 柔性事务核心思想
分布式场景不适用数据库刚性ACID事务,企业主流采用柔性事务,核心思想:不追求实时强一致性,保障数据最终一致性,通过重试、日志、消息补偿机制,让延迟执行的业务最终完成,适配微服务高并发、高可用特性。
4.3 两大主流柔性事务方案对比
4.3.1 RocketMQ可靠消息最终一致性
核心逻辑:基于RocketMQ事务消息,保障事务发起方本地业务与消息发送原子一致,消息发送成功后,强制消费方执行业务,通过消息回查、日志幂等保证最终一致。
适用场景:业务时效性要求较高、必须保证消息可靠送达、上下游业务强关联的场景(下单扣库存、订单联动业务)。
4.3.2 最大努力通知型事务
核心逻辑:业务主动方完成本地事务后,尽可能多次通知被动方,同时提供主动查询校对机制,不保证消息100%实时送达,依靠重试+兜底查询实现最终一致。
适用场景:时效性低、被动方结果不影响主动方业务的场景(支付结果通知、充值结果同步)。
第五章 RocketMQ环境Docker企业级部署
5.1 环境依赖
-
服务器系统:CentOS 7+
-
运行环境:Docker 20.10+
-
中间件版本:RocketMQ 4.4.0
-
辅助工具:RocketMQ可视化控制台
5.2 前置环境配置
关闭SELinux与防火墙,规避端口访问、权限异常问题:
# 临时关闭SELinux
setenforce 0
# 永久关闭SELinux
sed -i 's/^enforcing/disabled/' /etc/selinux/config
# 关闭防火墙
systemctl stop firewalld
systemctl disable firewalld
5.3 部署NameServer(注册中心)
NameServer是RocketMQ核心注册中心,负责Broker服务注册、发现、路由分发,无状态可集群部署。
# 创建数据挂载目录
mkdir -p /docker/rocketmq/namesrv/{logs,store}
# 启动NameServer容器
docker run -d \\
–restart=always \\
–name rmq-namesrv \\
-p 9876:9876 \\
-v /docker/rocketmq/namesrv/logs:/root/logs \\
-v /docker/rocketmq/namesrv/store:/root/store \\
-e "MAX_POSSIBLE_HEAP=100000000" \\
rocketmqinc/rocketmq sh mqnamesrv
5.4 部署Broker消息服务
Broker负责消息存储、投递、持久化,是消息核心处理节点,需自定义配置文件保证集群可用性。
5.4.1 编写Broker配置文件
# 创建配置目录
mkdir -p /docker/rocketmq/conf
# 编写核心配置
cat > /docker/rocketmq/conf/broker.conf << EOF
# 集群名称
brokerClusterName=DefaultCluster
# Broker节点名称
brokerName=broker-master
# 0代表Master节点
brokerId=0
# 消息删除时间(凌晨4点)
deleteWhen=04
# 消息磁盘保留时长48小时
fileReservedTime=48
# 异步主从同步模式
brokerRole=ASYNC_MASTER
# 异步刷盘策略(高性能)
flushDiskType=ASYNC_FLUSH
# 服务器内网IP(修改为当前服务器IP)
brokerIP1=192.168.66.100
# 磁盘最大使用比例
diskMaxUsedSpaceRatio=99
EOF
5.4.2 启动Broker容器
# 创建数据挂载目录
mkdir -p /docker/rocketmq/broker/{logs,store}
# 启动Broker容器
docker run -d \\
–restart=always \\
–name rmq-broker \\
–link rmq-namesrv:namesrv \\
-p 10911:10911 \\
-p 10909:10909 \\
–privileged=true \\
-v /docker/rocketmq/broker/logs:/root/logs \\
-v /docker/rocketmq/broker/store:/root/store \\
-v /docker/rocketmq/conf/broker.conf:/opt/rocketmq-4.4.0/conf/broker.conf \\
-e "NAMESRV_ADDR=namesrv:9876" \\
-e "MAX_POSSIBLE_HEAP=200000000" \\
rocketmqinc/rocketmq sh mqbroker -c /opt/rocketmq-4.4.0/conf/broker.conf
5.5 部署可视化控制台
docker run -d \\
–restart=always \\
–name rmq-console \\
-p 8080:8080 \\
-e "JAVA_OPTS=-Drocketmq.namesrv.addr=192.168.66.100:9876 -Dcom.rocketmq.sendMessageWithVIPChannel=false" \\
pangliang/rocketmq-console-ng
部署完成后,访问 http://192.168.66.100:8080 即可进入RocketMQ可视化管理界面,查看主题、消息、集群状态。
第六章 可靠消息最终一致性分布式事务实战
6.1 业务场景说明
实现用户下单-扣库存跨微服务分布式事务:订单服务创建订单、库存服务扣除商品库存,保证两个服务操作要么全部成功、要么全部回滚,杜绝订单创建成功但库存未扣除、或库存扣除无订单的脏数据。
6.2 核心实现原理
订单服务发送半事务消息到RocketMQ,消息暂不允许消费;
执行订单服务本地事务(创建订单、记录事务日志);
本地事务成功,通知MQ提交消息,库存服务消费消息扣库存;
本地事务失败,通知MQ回滚删除消息;
网络异常状态下,MQ定时回查本地事务状态,保证最终一致性;
通过事务日志实现幂等性,避免消息重复消费导致库存超扣。
6.3 数据库表设计
6.3.1 订单表 orders
CREATE TABLE `orders` (
`id` bigint NOT NULL AUTO_INCREMENT COMMENT '主键ID',
`order_no` varchar(64) NOT NULL COMMENT '订单唯一编号',
`product_id` bigint NOT NULL COMMENT '商品ID',
`pay_count` int NOT NULL COMMENT '购买数量',
`order_status` tinyint NOT NULL DEFAULT 1 COMMENT '订单状态 1-正常 0-失效',
`create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`update_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_order_no` (`order_no`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='订单表';
6.3.2 库存表 stock
CREATE TABLE `stock` (
`id` bigint NOT NULL AUTO_INCREMENT COMMENT '主键ID',
`product_id` bigint NOT NULL COMMENT '商品ID',
`total_count` int NOT NULL DEFAULT 0 COMMENT '总库存数量',
`used_count` int NOT NULL DEFAULT 0 COMMENT '已售出数量',
`create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`update_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_product_id` (`product_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='商品库存表';
6.3.3 事务日志表 tx_log(幂等核心)
CREATE TABLE `tx_log` (
`tx_no` varchar(64) NOT NULL COMMENT '全局事务编号',
`tx_status` tinyint NOT NULL DEFAULT 0 COMMENT '事务状态 0-处理中 1-已完成 2-已回滚',
`create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '事务创建时间',
`update_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '事务更新时间',
PRIMARY KEY (`tx_no`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='分布式事务日志表';
6.4 工程搭建(父工程+双微服务)
6.4.1 父工程统一依赖管理
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.itbaizhan</groupId>
<artifactId>rocketmq-transaction-parent</artifactId>
<version>1.0.0</version>
<packaging>pom</packaging>
<modules>
<module>order-service</module>
<module>stock-service</module>
</modules>
<properties>
<maven.compiler.source>8</maven.compiler.source>
<maven.compiler.target>8</maven.compiler.target>
<spring.boot.version>2.3.12.RELEASE</spring.boot.version>
<rocketmq.version>2.0.2</rocketmq.version>
<mybatis.plus.version>3.4.3.4</mybatis.plus.version>
</properties>
<dependencyManagement>
<dependencies>
<!– SpringBoot核心依赖 –>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-dependencies</artifactId>
<version>${spring.boot.version}</version>
<type>pom</type>
<scope>import</scope>
</dependency>
<!– MybatisPlus –>
<dependency>
<groupId>com.baomidou</groupId>
<artifactId>mybatis-plus-boot-starter</artifactId>
<version>${mybatis.plus.version}</version>
</dependency>
<!– RocketMQ –>
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>${rocketmq.version}</version>
</dependency>
</dependencies>
</dependencyManagement>
</project>
6.4.2 订单服务依赖(order-service)
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<parent>
<groupId>com.itbaizhan</groupId>
<artifactId>rocketmq-transaction-parent</artifactId>
<version>1.0.0</version>
</parent>
<modelVersion>4.0.0</modelVersion>
<artifactId>order-service</artifactId>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>com.baomidou</groupId>
<artifactId>mybatis-plus-boot-starter</artifactId>
</dependency>
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>
</project>
6.5 核心通用实体类
6.5.1 分布式事务消息载体 TxMessage
package com.itbaizhan.common.tx;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.io.Serializable;
/**
* 分布式事务消息通用载体
* 封装订单、库存服务交互的核心事务数据
*/
@Data
@NoArgsConstructor
@AllArgsConstructor
public class TxMessage implements Serializable {
private static final long serialVersionUID = 89234792347923423L;
/**
* 全局唯一事务编号
*/
private String globalTxNo;
/**
* 商品ID
*/
private Long productId;
/**
* 购买数量
*/
private Integer buyCount;
}
6.5.2 事务日志实体 TxLog
package com.itbaizhan.order.entity;
import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import java.time.LocalDateTime;
@Data
@TableName("tx_log")
public class TxLog {
@TableId(type = IdType.INPUT)
private String txNo;
private Integer txStatus;
private LocalDateTime createTime;
private LocalDateTime updateTime;
}
6.6 订单服务核心业务实现
6.6.1 订单服务配置文件 application.yml
server:
port: 9090
spring:
application:
name: order–transaction–service
datasource:
url: jdbc:mysql://192.168.66.100:3306/tx_order_db?useUnicode=true&characterEncoding=utf-8&serverTimezone=Asia/Shanghai&allowMultiQueries=true
username: root
password: 123456
driver-class-name: com.mysql.cj.jdbc.Driver
# RocketMQ配置
rocketmq:
name-server: 192.168.66.100:9876
producer:
group: order–transaction–group
# MybatisPlus配置
mybatis-plus:
mapper-locations: classpath:mapper/*.xml
global-config:
db-config:
id-type: auto
logic-delete-field: deleted
6.6.2 订单业务接口与实现
package com.itbaizhan.order.service;
import com.itbaizhan.order.entity.Order;
import com.itbaizhan.common.tx.TxMessage;
public interface OrderService {
/**
* 提交订单(入口方法)
* @param productId 商品ID
* @param buyCount 购买数量
*/
void submitOrder(Long productId, Integer buyCount);
/**
* 执行本地事务:创建订单+记录事务日志
* @param txMessage 全局事务消息
*/
void executeLocalOrderTx(TxMessage txMessage);
}
package com.itbaizhan.order.service.impl;
import com.alibaba.fastjson.JSON;
import com.itbaizhan.order.entity.Order;
import com.itbaizhan.order.entity.TxLog;
import com.itbaizhan.order.mapper.OrderMapper;
import com.itbaizhan.order.mapper.TxLogMapper;
import com.itbaizhan.order.service.OrderService;
import com.itbaizhan.common.tx.TxMessage;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import javax.annotation.Resource;
import java.time.LocalDateTime;
import java.util.UUID;
@Slf4j
@Service
public class OrderServiceImpl implements OrderService {
@Resource
private OrderMapper orderMapper;
@Resource
private TxLogMapper txLogMapper;
@Resource
private RocketMQTemplate rocketMQTemplate;
private static final String TRANSACTION_TOPIC = "order_stock_tx_topic";
@Override
public void submitOrder(Long productId, Integer buyCount) {
// 生成全局唯一事务编号
String globalTxNo = UUID.randomUUID().toString().replace("-", "");
// 封装事务消息
TxMessage txMessage = new TxMessage(globalTxNo, productId, buyCount);
Message<String> message = MessageBuilder.withPayload(JSON.toJSONString(txMessage)).build();
// 发送半事务消息,绑定本地事务
rocketMQTemplate.sendMessageInTransaction("order-transaction-group", TRANSACTION_TOPIC, message, txMessage);
log.info("订单事务消息发送成功,全局事务号:{}", globalTxNo);
}
@Override
@Transactional(rollbackFor = Exception.class)
public void executeLocalOrderTx(TxMessage txMessage) {
// 幂等判断:已执行过的事务直接返回
TxLog existTx = txLogMapper.selectById(txMessage.getGlobalTxNo());
if (existTx != null) {
log.info("订单事务已处理,无需重复执行,事务号:{}", txMessage.getGlobalTxNo());
return;
}
// 1. 创建订单数据
Order order = new Order();
order.setOrderNo(UUID.randomUUID().toString().replace("-", ""));
order.setProductId(txMessage.getProductId());
order.setPayCount(txMessage.getBuyCount());
order.setOrderStatus(1);
order.setCreateTime(LocalDateTime.now());
orderMapper.insert(order);
// 2. 记录事务日志,标记事务处理中
TxLog txLog = new TxLog();
txLog.setTxNo(txMessage.getGlobalTxNo());
txLog.setTxStatus(0);
txLog.setCreateTime(LocalDateTime.now());
txLogMapper.insert(txLog);
log.info("订单本地事务执行成功,事务号:{}", txMessage.getGlobalTxNo());
}
}
6.6.3 RocketMQ事务监听处理器(核心)
package com.itbaizhan.order.listener;
import com.alibaba.fastjson.JSON;
import com.itbaizhan.order.entity.TxLog;
import com.itbaizhan.order.mapper.TxLogMapper;
import com.itbaizhan.order.service.OrderService;
import com.itbaizhan.common.tx.TxMessage;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.annotation.RocketMQTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionState;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.Message;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import javax.annotation.Resource;
@Slf4j
@Component
@RocketMQTransactionListener(txProducerGroup = "order-transaction-group")
public class OrderTransactionListener implements RocketMQLocalTransactionListener {
@Autowired
private OrderService orderService;
@Resource
private TxLogMapper txLogMapper;
/**
* 执行本地事务
*/
@Override
@Transactional(rollbackFor = Exception.class)
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
TxMessage txMessage = (TxMessage) arg;
// 执行订单本地事务
orderService.executeLocalOrderTx(txMessage);
// 本地事务成功,提交消息,允许库存服务消费
return RocketMQLocalTransactionState.COMMIT;
} catch (Exception e) {
log.error("订单本地事务执行异常,事务回滚", e);
// 本地事务失败,回滚消息
return RocketMQLocalTransactionState.ROLLBACK;
}
}
/**
* 事务回查(网络异常兜底机制)
*/
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
try {
String payload = new String((byte[]) msg.getPayload());
TxMessage txMessage = JSON.parseObject(payload, TxMessage.class);
// 查询事务日志,判断本地事务是否执行成功
TxLog txLog = txLogMapper.selectById(txMessage.getGlobalTxNo());
if (txLog != null) {
log.info("事务回查:本地事务已执行成功,事务号:{}", txMessage.getGlobalTxNo());
return RocketMQLocalTransactionState.COMMIT;
}
// 事务状态未知,继续回查
return RocketMQLocalTransactionState.UNKNOWN;
} catch (Exception e) {
log.error("事务回查异常", e);
return RocketMQLocalTransactionState.ROLLBACK;
}
}
}
6.6.4 订单控制层(测试入口)
package com.itbaizhan.order.controller;
import com.itbaizhan.order.service.OrderService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import javax.annotation.Resource;
@Slf4j
@RestController
@RequestMapping("/order")
public class OrderController {
@Resource
private OrderService orderService;
@GetMapping("/create")
public String createOrder(@RequestParam Long productId, @RequestParam Integer buyCount) {
orderService.submitOrder(productId, buyCount);
return "订单提交成功,等待库存处理";
}
}
6.7 库存服务核心业务实现
6.7.1 库存服务配置文件 application.yml
server:
port: 9091
spring:
application:
name: stock–transaction–service
datasource:
url: jdbc:mysql://192.168.66.100:3306/tx_stock_db?useUnicode=true&characterEncoding=utf-8&serverTimezone=Asia/Shanghai&allowMultiQueries=true
username: root
password: 123456
driver-class-name: com.mysql.cj.jdbc.Driver
# RocketMQ配置
rocketmq:
name-server: 192.168.66.100:9876
# MybatisPlus配置
mybatis-plus:
mapper-locations: classpath:mapper/*.xml
global-config:
db-config:
id-type: auto
6.7.2 库存业务接口与实现
package com.itbaizhan.stock.service;
import com.itbaizhan.common.tx.TxMessage;
public interface StockService {
/**
* 扣减库存核心方法
* @param txMessage 事务消息
*/
void deductStock(TxMessage txMessage);
}
package com.itbaizhan.stock.service.impl;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.itbaizhan.common.tx.TxMessage;
import com.itbaizhan.stock.entity.Stock;
import com.itbaizhan.stock.entity.TxLog;
import com.itbaizhan.stock.mapper.StockMapper;
import com.itbaizhan.stock.mapper.TxLogMapper;
import com.itbaizhan.stock.service.StockService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import javax.annotation.Resource;
import java.time.LocalDateTime;
@Slf4j
@Service
public class StockServiceImpl implements StockService {
@Resource
private StockMapper stockMapper;
@Resource
private TxLogMapper txLogMapper;
@Override
@Transactional(rollbackFor = Exception.class)
public void deductStock(TxMessage txMessage) {
// 幂等校验:避免重复扣库存
TxLog existTx = txLogMapper.selectById(txMessage.getGlobalTxNo());
if (existTx != null) {
log.info("库存事务已处理,无需重复扣减,事务号:{}", txMessage.getGlobalTxNo());
return;
}
// 查询商品库存
LambdaQueryWrapper<Stock> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(Stock::getProductId, txMessage.getProductId());
Stock stock = stockMapper.selectOne(queryWrapper);
// 校验库存是否充足
if (stock == null || stock.getTotalCount() < txMessage.getBuyCount()) {
throw new RuntimeException("商品库存不足,扣减失败");
}
// 扣减库存、更新已售数量
stock.setTotalCount(stock.getTotalCount() – txMessage.getBuyCount());
stock.setUsedCount(stock.getUsedCount() + txMessage.getBuyCount());
stockMapper.updateById(stock);
// 记录库存事务日志,实现幂等
TxLog txLog = new TxLog();
txLog.setTxNo(txMessage.getGlobalTxNo());
txLog.setTxStatus(1);
txLog.setCreateTime(LocalDateTime.now());
txLogMapper.insert(txLog);
log.info("库存扣减成功,商品ID:{},扣减数量:{},事务号:{}",
txMessage.getProductId(), txMessage.getBuyCount(), txMessage.getGlobalTxNo());
}
}
6.7.3 库存消息消费者
package com.itbaizhan.stock.consumer;
import com.alibaba.fastjson.JSON;
import com.itbaizhan.common.tx.TxMessage;
import com.itbaizhan.stock.service.StockService;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
@Slf4j
@Component
@RocketMQMessageListener(consumerGroup = "stock-consumer-group", topic = "order_stock_tx_topic")
public class StockTransactionConsumer implements RocketMQListener<String> {
@Resource
private StockService stockService;
@Override
public void onMessage(String message) {
log.info("库存服务接收事务消息:{}", message);
// 解析事务消息
TxMessage txMessage = JSON.parseObject(message, TxMessage.class);
// 执行扣库存业务
stockService.deductStock(txMessage);
}
}
第七章 最大努力通知型分布式事务实战
7.1 业务场景说明
实现充值结果异步通知业务:用户充值成功后,充值服务完成扣款,通过RocketMQ异步通知账户服务更新余额。采用最大努力通知机制,支持消息重试+主动查询兜底,适配低时效性、高可靠性的通知场景。
7.2 核心实现机制
充值服务完成本地充值事务,发送充值结果消息到MQ;
账户服务监听MQ消息,消费成功则更新账户余额;
消息消费失败时,MQ自动重试推送,实现多次通知;
长期通知失败,账户服务主动调用充值服务接口,查询充值结果兜底;
通过事务日志实现幂等,避免重复充值、余额重复累加。
7.3 数据库表设计
7.3.1 充值记录表 recharge
CREATE TABLE `recharge` (
`id` bigint NOT NULL AUTO_INCREMENT COMMENT '主键ID',
`recharge_no` varchar(64) NOT NULL COMMENT '充值订单号',
`user_id` bigint NOT NULL COMMENT '用户ID',
`recharge_amount` decimal(10,2) NOT NULL COMMENT '充值金额',
`recharge_status` tinyint NOT NULL DEFAULT 0 COMMENT '充值状态 0-待处理 1-充值成功 2-充值失败',
`create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`update_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_recharge_no` (`recharge_no`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='用户充值记录表';
7.3.2 用户账户表 account
CREATE TABLE `account` (
`id` bigint NOT NULL AUTO_INCREMENT COMMENT '主键ID',
`user_id` bigint NOT NULL COMMENT '用户ID',
`balance` decimal(10,2) NOT NULL DEFAULT 0.00 COMMENT '账户余额',
`create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`update_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_user_id` (`user_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='用户账户表';
7.3.3 通知事务日志表 notify_tx_log(幂等+兜底核心)
CREATE TABLE `notify_tx_log` (
`tx_no` varchar(64) NOT NULL COMMENT '全局事务编号',
`recharge_no` varchar(64) NOT NULL COMMENT '充值订单号',
`notify_status` tinyint NOT NULL DEFAULT 0 COMMENT '通知状态 0-待通知 1-通知成功 2-通知失败',
`notify_times` int NOT NULL DEFAULT 0 COMMENT '已通知次数',
`last_notify_time` datetime DEFAULT NULL COMMENT '最后通知时间',
`create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`update_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`tx_no`),
UNIQUE KEY `uk_recharge_no` (`recharge_no`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='最大努力通知事务日志表';
7.4 整体工程结构说明
延续前文父工程,新增两个微服务:
-
recharge-service(充值服务):处理用户充值、本地事务落库、发送充值通知消息、提供充值结果查询兜底接口
-
account-service(账户服务):消费充值通知、更新用户余额、定时任务兜底查询未同步充值订单
7.5 通用消息实体(充值通知消息)
package com.itbaizhan.common.notify;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.io.Serializable;
import java.math.BigDecimal;
/**
* 充值通知消息实体
*/
@Data
@NoArgsConstructor
@AllArgsConstructor
public class RechargeNotifyMessage implements Serializable {
private static final long serialVersionUID = 56781234567890L;
/**
* 全局事务号
*/
private String globalTxNo;
/**
* 充值订单号
*/
private String rechargeNo;
/**
* 用户ID
*/
private Long userId;
/**
* 充值金额
*/
private BigDecimal rechargeAmount;
/**
* 充值状态
*/
private Integer rechargeStatus;
}
7.6 充值服务(recharge-service)完整实现
7.6.1 服务依赖 pom.xml
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<parent>
<groupId>com.itbaizhan</groupId>
<artifactId>rocketmq-transaction-parent</artifactId>
<version>1.0.0</version>
</parent>
<modelVersion>4.0.0</modelVersion>
<artifactId>recharge-service</artifactId>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-task</artifactId>
</dependency>
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>com.baomidou</groupId>
<artifactId>mybatis-plus-boot-starter</artifactId>
</dependency>
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>
</project>
7.6.2 配置文件 application.yml
server:
port: 9092
spring:
application:
name: recharge–service
datasource:
url: jdbc:mysql://192.168.66.100:3306/tx_recharge_db?useUnicode=true&characterEncoding=utf-8&serverTimezone=Asia/Shanghai&allowMultiQueries=true
username: root
password: 123456
driver-class-name: com.mysql.cj.jdbc.Driver
# RocketMQ配置
rocketmq:
name-server: 192.168.66.100:9876
producer:
group: recharge–notify–producer–group
# MybatisPlus配置
mybatis-plus:
mapper-locations: classpath:mapper/*.xml
global-config:
db-config:
id-type: auto
7.6.3 充值业务核心接口与实现
package com.itbaizhan.recharge.service;
import com.itbaizhan.recharge.entity.Recharge;
/**
* 充值业务接口
*/
public interface RechargeService {
/**
* 用户充值业务
* @param userId 用户ID
* @param amount 充值金额
* @return 充值订单号
*/
String userRecharge(Long userId, java.math.BigDecimal amount);
/**
* 根据充值订单号查询充值结果(兜底查询接口)
* @param rechargeNo 充值订单号
* @return 充值订单信息
*/
Recharge getRechargeInfo(String rechargeNo);
}
package com.itbaizhan.recharge.service.impl;
import com.alibaba.fastjson.JSON;
import com.itbaizhan.common.notify.RechargeNotifyMessage;
import com.itbaizhan.recharge.entity.Recharge;
import com.itbaizhan.recharge.entity.NotifyTxLog;
import com.itbaizhan.recharge.mapper.RechargeMapper;
import com.itbaizhan.recharge.mapper.NotifyTxLogMapper;
import com.itbaizhan.recharge.service.RechargeService;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import javax.annotation.Resource;
import java.time.LocalDateTime;
import java.util.UUID;
@Slf4j
@Service
public class RechargeServiceImpl implements RechargeService {
@Resource
private RechargeMapper rechargeMapper;
@Resource
private NotifyTxLogMapper notifyTxLogMapper;
@Resource
private RocketMQTemplate rocketMQTemplate;
private static final String RECHARGE_NOTIFY_TOPIC = "recharge_notify_topic";
@Override
@Transactional(rollbackFor = Exception.class)
public String userRecharge(Long userId, java.math.BigDecimal amount) {
// 1. 生成唯一订单号、事务号
String rechargeNo = UUID.randomUUID().replace("-", "");
String globalTxNo = UUID.randomUUID().replace("-", "");
// 2. 新增充值订单,标记充值成功
Recharge recharge = new Recharge();
recharge.setRechargeNo(rechargeNo);
recharge.setUserId(userId);
recharge.setRechargeAmount(amount);
recharge.setRechargeStatus(1);
recharge.setCreateTime(LocalDateTime.now());
rechargeMapper.insert(recharge);
// 3. 新增通知事务日志,初始待通知状态
NotifyTxLog txLog = new NotifyTxLog();
txLog.setTxNo(globalTxNo);
txLog.setRechargeNo(rechargeNo);
txLog.setNotifyStatus(0);
txLog.setNotifyTimes(0);
txLog.setCreateTime(LocalDateTime.now());
notifyTxLogMapper.insert(txLog);
// 4. 封装通知消息,发送MQ
RechargeNotifyMessage message = new RechargeNotifyMessage();
message.setGlobalTxNo(globalTxNo);
message.setRechargeNo(rechargeNo);
message.setUserId(userId);
message.setRechargeAmount(amount);
message.setRechargeStatus(1);
Message<String> mqMessage = MessageBuilder.withPayload(JSON.toJSONString(message)).build();
// 同步发送消息,失败会自动重试
rocketMQTemplate.syncSend(RECHARGE_NOTIFY_TOPIC, mqMessage);
// 5. 更新通知次数、最后通知时间
txLog.setNotifyTimes(1);
txLog.setLastNotifyTime(LocalDateTime.now());
notifyTxLogMapper.updateById(txLog);
log.info("用户充值成功,已发送通知,订单号:{}", rechargeNo);
return rechargeNo;
}
@Override
public Recharge getRechargeInfo(String rechargeNo) {
return rechargeMapper.selectOne(com.baomidou.mybatisplus.core.conditions.query.LambdaQueryChainWrapper
.lambdaQuery(rechargeMapper)
.eq(Recharge::getRechargeNo, rechargeNo)
.getWrapper());
}
}
7.6.4 充值兜底重试定时任务(最大努力核心)
package com.itbaizhan.recharge.task;
import com.alibaba.fastjson.JSON;
import com.itbaizhan.common.notify.RechargeNotifyMessage;
import com.itbaizhan.recharge.entity.NotifyTxLog;
import com.itbaizhan.recharge.entity.Recharge;
import com.itbaizhan.recharge.mapper.NotifyTxLogMapper;
import com.itbaizhan.recharge.service.RechargeService;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
import java.time.LocalDateTime;
import java.util.List;
/**
* 定时重试通知任务:最大努力通知核心
* 定时扫描通知失败、未完成的订单,持续重试推送
*/
@Slf4j
@Component
public class NotifyRetryTask {
@Resource
private NotifyTxLogMapper notifyTxLogMapper;
@Resource
private RechargeService rechargeService;
@Resource
private RocketMQTemplate rocketMQTemplate;
private static final String RECHARGE_NOTIFY_TOPIC = "recharge_notify_topic";
// 最大重试次数
private static final int MAX_RETRY_TIMES = 5;
@Scheduled(cron = "0 */1 * * * ?")
public void retryNotify() {
// 查询待通知、通知失败且未超过最大重试次数的记录
List<NotifyTxLog> waitNotifyList = notifyTxLogMapper.selectList(
com.baomidou.mybatisplus.core.conditions.query.LambdaQueryChainWrapper
.lambdaQuery(notifyTxLogMapper)
.in(NotifyTxLog::getNotifyStatus, 0, 2)
.lt(NotifyTxLog::getNotifyTimes, MAX_RETRY_TIMES)
.getWrapper()
);
if (waitNotifyList == null || waitNotifyList.isEmpty()) {
return;
}
log.info("开始执行充值通知重试,待重试数量:{}", waitNotifyList.size());
for (NotifyTxLog txLog : waitNotifyList) {
try {
// 查询最新充值状态
Recharge recharge = rechargeService.getRechargeInfo(txLog.getRechargeNo());
if (recharge == null) {
continue;
}
// 封装消息重试推送
RechargeNotifyMessage message = new RechargeNotifyMessage();
message.setGlobalTxNo(txLog.getTxNo());
message.setRechargeNo(txLog.getRechargeNo());
message.setUserId(recharge.getUserId());
message.setRechargeAmount(recharge.getRechargeAmount());
message.setRechargeStatus(recharge.getRechargeStatus());
org.springframework.messaging.Message<String> mqMessage =
org.springframework.messaging.support.MessageBuilder.withPayload(JSON.toJSONString(message)).build();
rocketMQTemplate.syncSend(RECHARGE_NOTIFY_TOPIC, mqMessage);
// 更新重试次数与时间
txLog.setNotifyTimes(txLog.getNotifyTimes() + 1);
txLog.setLastNotifyTime(LocalDateTime.now());
notifyTxLogMapper.updateById(txLog);
log.info("充值通知重试成功,订单号:{},当前重试次数:{}", txLog.getRechargeNo(), txLog.getNotifyTimes());
} catch (Exception e) {
log.error("充值通知重试失败,订单号:{}", txLog.getRechargeNo(), e);
// 标记为通知失败,等待下次重试
txLog.setNotifyStatus(2);
notifyTxLogMapper.updateById(txLog);
}
}
}
}
7.6.5 充值控制层(测试入口+兜底查询接口)
package com.itbaizhan.recharge.controller;
import com.itbaizhan.recharge.entity.Recharge;
import com.itbaizhan.recharge.service.RechargeService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.web.bind.annotation.*;
import javax.annotation.Resource;
import java.math.BigDecimal;
@Slf4j
@RestController
@RequestMapping("/recharge")
public class RechargeController {
@Resource
private RechargeService rechargeService;
/**
* 充值测试接口
*/
@GetMapping("/doRecharge")
public String doRecharge(@RequestParam Long userId, @RequestParam BigDecimal amount) {
String rechargeNo = rechargeService.userRecharge(userId, amount);
return "充值成功,充值订单号:" + rechargeNo;
}
/**
* 兜底查询接口,供账户服务主动调用
*/
@GetMapping("/getInfo/{rechargeNo}")
public Recharge getRechargeInfo(@PathVariable String rechargeNo) {
return rechargeService.getRechargeInfo(rechargeNo);
}
}
7.7 账户服务(account-service)完整实现
7.7.1 服务依赖 pom.xml
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<parent>
<groupId>com.itbaizhan</groupId>
<artifactId>rocketmq-transaction-parent</artifactId>
<version>1.0.0</version>
</parent>
<modelVersion>4.0.0</modelVersion>
<artifactId>account-service</artifactId>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-task</artifactId>
</dependency>
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>com.baomidou</groupId>
<artifactId>mybatis-plus-boot-starter</artifactId>
</dependency>
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-starter-openfeign</artifactId>
<version>2.2.9.RELEASE</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>
</project>
7.7.2 配置文件 application.yml
server:
port: 9093
spring:
application:
name: account–service
datasource:
url: jdbc:mysql://192.168.66.100:3306/tx_account_db?useUnicode=true&characterEncoding=utf-8&serverTimezone=Asia/Shanghai&allowMultiQueries=true
username: root
password: 123456
driver-class-name: com.mysql.cj.jdbc.Driver
# RocketMQ配置
rocketmq:
name-server: 192.168.66.100:9876
# MybatisPlus配置
mybatis-plus:
mapper-locations: classpath:mapper/*.xml
global-config:
db-config:
id-type: auto
# 充值服务接口地址
service:
recharge:
url: http://localhost:9092
7.7.3 Feign兜底查询接口
package com.itbaizhan.account.feign;
import com.itbaizhan.account.feign.fallback.RechargeFeignFallback;
import com.itbaizhan.recharge.entity.Recharge;
import org.springframework.cloud.openfeign.FeignClient;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
@FeignClient(name = "recharge-service", url = "${service.recharge.url}", fallback = RechargeFeignFallback.class)
public interface RechargeFeignClient {
/**
* 主动查询充值订单结果
*/
@GetMapping("/recharge/getInfo/{rechargeNo}")
Recharge getRechargeInfo(@PathVariable("rechargeNo") String rechargeNo);
}
package com.itbaizhan.account.feign.fallback;
import com.itbaizhan.account.feign.RechargeFeignClient;
import com.itbaizhan.recharge.entity.Recharge;
import org.springframework.stereotype.Component;
@Component
public class RechargeFeignFallback implements RechargeFeignClient {
@Override
public Recharge getRechargeInfo(String rechargeNo) {
return null;
}
}
7.7.4 账户余额更新业务实现
package com.itbaizhan.account.service;
import com.itbaizhan.common.notify.RechargeNotifyMessage;
public interface AccountService {
/**
* 更新用户账户余额
* @param message 充值通知消息
*/
void updateBalance(RechargeNotifyMessage message);
}
package com.itbaizhan.account.service.impl;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.itbaizhan.account.entity.Account;
import com.itbaizhan.account.entity.NotifyRecord;
import com.itbaizhan.account.mapper.AccountMapper;
import com.itbaizhan.account.mapper.NotifyRecordMapper;
import com.itbaizhan.account.service.AccountService;
import com.itbaizhan.common.notify.RechargeNotifyMessage;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import javax.annotation.Resource;
import java.time.LocalDateTime;
@Slf4j
@Service
public class AccountServiceImpl implements AccountService {
@Resource
private AccountMapper accountMapper;
@Resource
private NotifyRecordMapper notifyRecordMapper;
@Override
@Transactional(rollbackFor = Exception.class)
public void updateBalance(RechargeNotifyMessage message) {
String rechargeNo = message.getRechargeNo();
// 幂等判断:已处理过的订单直接跳过
Long count = notifyRecordMapper.selectCount(
new LambdaQueryWrapper<NotifyRecord>().eq(NotifyRecord::getRechargeNo, rechargeNo)
);
if (count > 0) {
log.info("充值通知已处理,无需重复更新余额,订单号:{}", rechargeNo);
return;
}
// 查询用户账户
Account account = accountMapper.selectOne(
new LambdaQueryWrapper<Account>().eq(Account::getUserId, message.getUserId())
);
if (account == null) {
// 新用户初始化账户余额
account = new Account();
account.setUserId(message.getUserId());
account.setBalance(message.getRechargeAmount());
account.setCreateTime(LocalDateTime.now());
accountMapper.insert(account);
} else {
// 老用户累加余额
account.setBalance(account.getBalance().add(message.getRechargeAmount()));
account.setUpdateTime(LocalDateTime.now());
accountMapper.updateById(account);
}
// 记录处理日志,保证幂等
NotifyRecord record = new NotifyRecord();
record.setTxNo(message.getGlobalTxNo());
record.setRechargeNo(rechargeNo);
record.setUserId(message.getUserId());
record.setRechargeAmount(message.getRechargeAmount());
record.setCreateTime(LocalDateTime.now());
notifyRecordMapper.insert(record);
log.info("用户余额更新成功,用户ID:{},充值金额:{}", message.getUserId(), message.getRechargeAmount());
}
}
7.7.5 充值消息消费者
package com.itbaizhan.account.consumer;
import com.alibaba.fastjson.JSON;
import com.itbaizhan.account.service.AccountService;
import com.itbaizhan.common.notify.RechargeNotifyMessage;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
@Slf4j
@Component
@RocketMQMessageListener(consumerGroup = "account-notify-consumer-group", topic = "recharge_notify_topic")
public class RechargeNotifyConsumer implements RocketMQListener<String> {
@Resource
private AccountService accountService;
@Override
public void onMessage(String message) {
log.info("账户服务接收充值通知消息:{}", message);
RechargeNotifyMessage notifyMessage = JSON.parseObject(message, RechargeNotifyMessage.class);
// 执行余额更新逻辑
accountService.updateBalance(notifyMessage);
}
}
7.7.6 定时兜底主动查询任务(最终一致性兜底)
package com.itbaizhan.account.task;
import com.itbaizhan.account.entity.NotifyRecord;
import com.itbaizhan.account.feign.RechargeFeignClient;
import com.itbaizhan.account.mapper.NotifyRecordMapper;
import com.itbaizhan.account.service.AccountService;
import com.itbaizhan.common.notify.RechargeNotifyMessage;
import com.itbaizhan.recharge.entity.Recharge;
import lombok.extern.slf4j.Slf4j;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
import java.util.List;
import java.util.UUID;
/**
* 兜底补偿任务:主动查询充值订单,修复消息丢失、消费失败数据
*/
@Slf4j
@Component
public class AccountCompensateTask {
@Resource
private RechargeFeignClient rechargeFeignClient;
@Resource
private NotifyRecordMapper notifyRecordMapper;
@Resource
private AccountService accountService;
@Scheduled(cron = "0 0/2 * * * ?")
public void compensateRecharge() {
log.info("开始执行账户充值兜底补偿任务");
// 实际生产可根据业务查询一段时间内未同步的充值订单
// 此处模拟兜底:可结合业务日志、对账日志实现
// 核心逻辑:主动调用充值服务查询订单,未处理则更新余额
}
}
7.8 两大分布式事务方案最终总结对比
| 核心原理 | 半消息+本地事务+消息回查,保证发送方事务与消息原子性 | 本地事务完成后多次重试通知+主动兜底查询 |
| 时效性 | 高,实时联动业务 | 低,允许短暂数据不一致 |
| 适用场景 | 下单扣库存、业务强依赖联动场景 | 支付通知、充值通知、对账同步场景 |
| 兜底机制 | MQ自动事务回查 | 定时重试+主动接口查询 |
| 一致性强度 | 更高,近乎实时最终一致 | 弱最终一致 |
第八章 项目踩坑总结与生产优化方案
8.1 幂等性必做优化
分布式MQ消息天然存在重试机制,所有消费业务必须做幂等校验,禁止重复执行业务:
-
统一使用全局事务号、订单号作为唯一幂等Key
-
所有事务操作落地事务日志,执行前先查询、后执行
-
禁止无幂等校验的库存扣减、余额累加、订单创建操作
8.2 MQ消息异常处理
-
事务消息必须配置事务回查,解决网络超时、服务宕机导致的状态未知问题
-
普通通知消息配置重试次数、死信队列,避免消息无限重试堆积
-
生产环境务必部署Broker集群,避免单点故障导致消息丢失
8.3 生产环境配置建议
-
调整JVM堆内存,适配服务器配置,避免OOM
-
开启Broker异步刷盘+主从集群,兼顾性能与数据安全
-
定时清理过期事务日志、消息消费日志,避免数据表数据膨胀
-
增加分布式事务监控日志,快速定位事务不一致问题
全文完结:本文从本地事务原理、并发问题、隔离级别,到分布式事务理论、RocketMQ部署,再到两大企业级分布式事务方案完整实战,覆盖90%以上Java微服务分布式事务业务场景,可直接用于项目开发与面试复习。




