欢迎光临
我们一直在努力

Java高级全套教程(十二)—— 分布式事务超详细实战全解(基础理论+双方案企业级实战)

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: ordertransactionservice
    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: ordertransactiongroup

    # 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: stocktransactionservice
    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: rechargeservice
    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: rechargenotifyproducergroup

    # 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&gt;
    &lt;/plugins&gt;
    &lt;/build&gt;
    &lt;/project&gt;

    7.7.2 配置文件 application.yml

    server:
    port: 9093

    spring:
    application:
    name: accountservice
    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 两大分布式事务方案最终总结对比

    对比维度RocketMQ可靠消息最终一致性最大努力通知型
    核心原理 半消息+本地事务+消息回查,保证发送方事务与消息原子性 本地事务完成后多次重试通知+主动兜底查询
    时效性 高,实时联动业务 低,允许短暂数据不一致
    适用场景 下单扣库存、业务强依赖联动场景 支付通知、充值通知、对账同步场景
    兜底机制 MQ自动事务回查 定时重试+主动接口查询
    一致性强度 更高,近乎实时最终一致 弱最终一致

    第八章 项目踩坑总结与生产优化方案

    8.1 幂等性必做优化

    分布式MQ消息天然存在重试机制,所有消费业务必须做幂等校验,禁止重复执行业务:

    • 统一使用全局事务号、订单号作为唯一幂等Key

    • 所有事务操作落地事务日志,执行前先查询、后执行

    • 禁止无幂等校验的库存扣减、余额累加、订单创建操作

    8.2 MQ消息异常处理

    • 事务消息必须配置事务回查,解决网络超时、服务宕机导致的状态未知问题

    • 普通通知消息配置重试次数、死信队列,避免消息无限重试堆积

    • 生产环境务必部署Broker集群,避免单点故障导致消息丢失

    8.3 生产环境配置建议

    • 调整JVM堆内存,适配服务器配置,避免OOM

    • 开启Broker异步刷盘+主从集群,兼顾性能与数据安全

    • 定时清理过期事务日志、消息消费日志,避免数据表数据膨胀

    • 增加分布式事务监控日志,快速定位事务不一致问题

    全文完结:本文从本地事务原理、并发问题、隔离级别,到分布式事务理论、RocketMQ部署,再到两大企业级分布式事务方案完整实战,覆盖90%以上Java微服务分布式事务业务场景,可直接用于项目开发与面试复习。

    赞(0)
    未经允许不得转载:171主机测评 » Java高级全套教程(十二)—— 分布式事务超详细实战全解(基础理论+双方案企业级实战)
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址