一、业务说明

1、分布式事务业务说明

业务说明
这里我们会创建三个服务,一个订单服务,一个库存服务,一个账户服务。
当用户下单时,会在订单服务中创建一个订单,然后通过远程调用库存服务来扣减下单商品的库存,再通过远程调用账户服务来扣减用户账户里面的余额,最后在订单服务中修改订单状态为已完成。
该操作跨越三个数据库,有两次远程调用,很明显会有分布式事务问题。

二、数据库创建

下订单–>扣库存–>减账户(余额)
1、创建业务数据库

seata_order: 存储订单的数据库
seata_storage:存储库存的数据库
seata_account: 存储账户信息的数据库

建表SQL

CREATE DATABASE seata_order;
CREATE DATABASE seata_storage;
CREATE DATABASE seata_account;

2、按照上述3库分别建对应业务表
按照上述3库分别建对应业务表
seata_order库下建t_order表
seata_storage库下建t_storage表
seata_account库下建t_account表

CREATE TABLE t_order(
    `id` BIGINT(11) NOT NULL AUTO_INCREMENT PRIMARY KEY,
    `user_id` BIGINT(11) DEFAULT NULL COMMENT '用户id',
    `product_id` BIGINT(11) DEFAULT NULL COMMENT '产品id',
    `count` INT(11) DEFAULT NULL COMMENT '数量',
    `money` DECIMAL(11,0) DEFAULT NULL COMMENT '金额',
    `status` INT(1) DEFAULT NULL COMMENT '订单状态:0:创建中; 1:已完结'
) ENGINE=INNODB AUTO_INCREMENT=7 DEFAULT CHARSET=utf8;
 
SELECT * FROM t_order;

CREATE TABLE t_storage(
    `id` BIGINT(11) NOT NULL AUTO_INCREMENT PRIMARY KEY,
    `product_id` BIGINT(11) DEFAULT NULL COMMENT '产品id',
    `total` INT(11) DEFAULT NULL COMMENT '总库存',
    `used` INT(11) DEFAULT NULL COMMENT '已用库存',
    `residue` INT(11) DEFAULT NULL COMMENT '剩余库存'
) ENGINE=INNODB AUTO_INCREMENT=2 DEFAULT CHARSET=utf8;
 
INSERT INTO seata_storage.t_storage(`id`,`product_id`,`total`,`used`,`residue`)VALUES('1','1','100','0','100');
 
SELECT * FROM t_storage;

CREATE TABLE t_account(
    `id` BIGINT(11) NOT NULL AUTO_INCREMENT PRIMARY KEY COMMENT 'id',
    `user_id` BIGINT(11) DEFAULT NULL COMMENT '用户id',
    `total` DECIMAL(10,0) DEFAULT NULL COMMENT '总额度',
    `used` DECIMAL(10,0) DEFAULT NULL COMMENT '已用余额',
    `residue` DECIMAL(10,0) DEFAULT '0' COMMENT '剩余可用额度'
) ENGINE=INNODB AUTO_INCREMENT=2 DEFAULT CHARSET=utf8;
 
INSERT INTO seata_account.t_account(`id`,`user_id`,`total`,`used`,`residue`) VALUES('1','1','1000','0','1000');
  
SELECT * FROM t_account;

3、3个库分别建对应的回滚日志表
按照上述3库分别建对应的回滚日志表
订单-库存-账户3个库下都需要建各自的回滚日志表seata\conf目录下的db_undo_log.sql

-- 注意:需要在 seata_order、seata_storage、seata_account 三个库中都执行此脚本
-- 事务回滚日志表:用于Seata AT模式中存储事务修改记录,实现分布式事务的回滚机制
-- 需在所有参与分布式事务的业务数据库(seata_order、seata_storage、seata_account)中创建
CREATE TABLE `undo_log` (
  `id` bigint(20) NOT NULL AUTO_INCREMENT COMMENT '自增主键,唯一标识一条回滚日志记录',
  `branch_id` bigint(20) NOT NULL COMMENT '分支事务ID,与全局事务关联的本地事务标识',
  `xid` varchar(100) NOT NULL COMMENT '全局事务ID,跨服务的分布式事务唯一标识',
  `context` varchar(128) NOT NULL COMMENT '上下文信息,存储事务相关的附加数据(JSON格式)',
  `rollback_info` longblob NOT NULL COMMENT '回滚信息,二进制存储的数据修改前后的镜像及SQL类型等关键信息',
  `log_status` int(11) NOT NULL COMMENT '日志状态:0-未提交,1-已提交,用于标记日志是否有效',
  `log_created` datetime NOT NULL COMMENT '日志创建时间,记录事务操作发生的时间',
  `log_modified` datetime NOT NULL COMMENT '日志修改时间,记录日志最后更新的时间',
  PRIMARY KEY (`id`),
  UNIQUE KEY `ux_undo_log` (`xid`,`branch_id`) COMMENT '联合唯一索引,确保同一全局事务下的分支事务日志唯一'
) ENGINE=InnoDB AUTO_INCREMENT=1 DEFAULT CHARSET=utf8 COMMENT='Seata AT模式事务回滚日志表,用于存储数据修改记录以支持事务回滚';

最终效果
在这里插入图片描述

三、3个项目构建

Seata 分布式事务示例:三个微服务项目实现

下面将创建三个 Spring Boot 微服务项目(订单服务、库存服务、账户服务),演示并通过 Seata 实现分布式事务管理。

1. 公共依赖与配置

所有项目都需要添加的核心依赖(pom.xml):

<dependencies>
    <!-- Spring Boot 核心 -->
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-web</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-data-jpa</artifactId>
    </dependency>
    <dependency>
        <groupId>mysql</groupId>
        <artifactId>mysql-connector-java</artifactId>
        <scope>runtime</scope>
    </dependency>
    
    <!-- Seata 依赖 -->
    <dependency>
        <groupId>io.seata</groupId>
        <artifactId>seata-spring-boot-starter</artifactId>
        <version>1.7.0</version>
    </dependency>
    <dependency>
        <groupId>com.alibaba.cloud</groupId>
        <artifactId>spring-cloud-starter-alibaba-nacos-discovery</artifactId>
        <version>2021.0.5.0</version>
    </dependency>
</dependencies>
2. 订单服务 (seata-order-service)
配置文件 (application.yml)
server:
  port: 8081

spring:
  application:
    name: seata-order-service
  datasource:
    driver-class-name: com.mysql.cj.jdbc.Driver
    url: jdbc:mysql://localhost:3306/seata_order?useUnicode=true&characterEncoding=utf8&useSSL=false
    username: root
    password: 123456
  cloud:
    nacos:
      discovery:
        server-addr: localhost:8848

seata:
  enabled: true
  application-id: ${spring.application.name}
  tx-service-group: default_tx_group  # 与 Seata 服务端配置一致
  registry:
    type: nacos
    nacos:
      server-addr: localhost:8848
      group: SEATA_GROUP
  service:
    vgroup-mapping:
      default_tx_group: default  # 与 Seata 服务端配置一致
核心代码

订单实体类:

@Entity
@Table(name = "t_order")
public class Order {
    @Id
    @GeneratedValue(strategy = GenerationType.IDENTITY)
    private Long id;
    private Long userId;
    private Long productId;
    private Integer count;
    private BigDecimal money;
    private Integer status;  // 0:创建中 1:已完成
    
    // getter 和 setter 省略
}

订单Repository:

public interface OrderRepository extends JpaRepository<Order, Long> {
}

远程服务调用接口:

@FeignClient(name = "seata-storage-service")
public interface StorageService {
    @PostMapping("/storage/decrease")
    Result decrease(@RequestParam("productId") Long productId, @RequestParam("count") Integer count);
}

@FeignClient(name = "seata-account-service")
public interface AccountService {
    @PostMapping("/account/decrease")
    Result decrease(@RequestParam("userId") Long userId, @RequestParam("money") BigDecimal money);
}

订单服务实现:

@Service
public class OrderService {
    @Autowired
    private OrderRepository orderRepository;
    @Autowired
    private StorageService storageService;
    @Autowired
    private AccountService accountService;

    /**
     * 创建订单,@GlobalTransactional 注解标记分布式事务
     */
    @GlobalTransactional(name = "create-order", rollbackFor = Exception.class)
    public Result create(Order order) {
        System.out.println("开始创建订单...");
        
        // 1. 创建订单
        order.setStatus(0);
        orderRepository.save(order);
        
        // 2. 扣减库存
        storageService.decrease(order.getProductId(), order.getCount());
        
        // 3. 扣减账户余额
        accountService.decrease(order.getUserId(), order.getMoney());
        
        // 4. 修改订单状态为已完成
        order.setStatus(1);
        orderRepository.save(order);
        
        System.out.println("订单创建完成");
        return Result.success(order);
    }
}

订单控制器:

@RestController
@RequestMapping("/order")
public class OrderController {
    @Autowired
    private OrderService orderService;
    
    @PostMapping("/create")
    public Result createOrder(Order order) {
        return orderService.create(order);
    }
}
3. 库存服务 (seata-storage-service)
配置文件 (application.yml)
server:
  port: 8082

spring:
  application:
    name: seata-storage-service
  datasource:
    driver-class-name: com.mysql.cj.jdbc.Driver
    url: jdbc:mysql://localhost:3306/seata_storage?useUnicode=true&characterEncoding=utf8&useSSL=false
    username: root
    password: 123456
  cloud:
    nacos:
      discovery:
        server-addr: localhost:8848

seata:
  enabled: true
  application-id: ${spring.application.name}
  tx-service-group: default_tx_group
  registry:
    type: nacos
    nacos:
      server-addr: localhost:8848
      group: SEATA_GROUP
  service:
    vgroup-mapping:
      default_tx_group: default
核心代码

库存实体类:

@Entity
@Table(name = "t_storage")
public class Storage {
    @Id
    @GeneratedValue(strategy = GenerationType.IDENTITY)
    private Long id;
    private Long productId;
    private Integer total;
    private Integer used;
    private Integer residue;
    
    // getter 和 setter 省略
}

库存服务实现:

@Service
public class StorageService {
    @Autowired
    private StorageRepository storageRepository;

    /**
     * 扣减库存
     */
    public Result decrease(Long productId, Integer count) {
        System.out.println("开始扣减库存...");
        
        Optional<Storage> storageOptional = storageRepository.findByProductId(productId);
        if (!storageOptional.isPresent()) {
            throw new RuntimeException("库存不存在");
        }
        
        Storage storage = storageOptional.get();
        if (storage.getResidue() < count) {
            throw new RuntimeException("库存不足");
        }
        
        // 扣减库存
        storage.setUsed(storage.getUsed() + count);
        storage.setResidue(storage.getResidue() - count);
        storageRepository.save(storage);
        
        System.out.println("库存扣减完成");
        return Result.success();
    }
}

库存控制器:

@RestController
@RequestMapping("/storage")
public class StorageController {
    @Autowired
    private StorageService storageService;
    
    @PostMapping("/decrease")
    public Result decrease(@RequestParam("productId") Long productId, @RequestParam("count") Integer count) {
        return storageService.decrease(productId, count);
    }
}
4. 账户服务 (seata-account-service)
配置文件 (application.yml)
server:
  port: 8083

spring:
  application:
    name: seata-account-service
  datasource:
    driver-class-name: com.mysql.cj.jdbc.Driver
    url: jdbc:mysql://localhost:3306/seata_account?useUnicode=true&characterEncoding=utf8&useSSL=false
    username: root
    password: 123456
  cloud:
    nacos:
      discovery:
        server-addr: localhost:8848

seata:
  enabled: true
  application-id: ${spring.application.name}
  tx-service-group: default_tx_group
  registry:
    type: nacos
    nacos:
      server-addr: localhost:8848
      group: SEATA_GROUP
  service:
    vgroup-mapping:
      default_tx_group: default
核心代码

账户实体类:

@Entity
@Table(name = "t_account")
public class Account {
    @Id
    @GeneratedValue(strategy = GenerationType.IDENTITY)
    private Long id;
    private Long userId;
    private BigDecimal total;
    private BigDecimal used;
    private BigDecimal residue;
    
    // getter 和 setter 省略
}

账户服务实现:

@Service
public class AccountService {
    @Autowired
    private AccountRepository accountRepository;

    /**
     * 扣减账户余额
     */
    public Result decrease(Long userId, BigDecimal money) {
        System.out.println("开始扣减账户余额...");
        
        // 模拟超时场景,测试事务回滚
        // try { Thread.sleep(20000); } catch (InterruptedException e) { e.printStackTrace(); }
        
        Optional<Account> accountOptional = accountRepository.findByUserId(userId);
        if (!accountOptional.isPresent()) {
            throw new RuntimeException("账户不存在");
        }
        
        Account account = accountOptional.get();
        if (account.getResidue().compareTo(money) < 0) {
            throw new RuntimeException("余额不足");
        }
        
        // 扣减余额
        account.setUsed(account.getUsed().add(money));
        account.setResidue(account.getResidue().subtract(money));
        accountRepository.save(account);
        
        System.out.println("账户余额扣减完成");
        return Result.success();
    }
}

账户控制器:

@RestController
@RequestMapping("/account")
public class AccountController {
    @Autowired
    private AccountService accountService;
    
    @PostMapping("/decrease")
    public Result decrease(@RequestParam("userId") Long userId, @RequestParam("money") BigDecimal money) {
        return accountService.decrease(userId, money);
    }
}
5. 测试分布式事务
测试步骤:
  1. 启动 Nacos 服务(默认端口 8848)
  2. 启动 Seata Server(配置已完成)
  3. 分别启动三个微服务:订单服务(8081)、库存服务(8082)、账户服务(8083)
  4. 发送创建订单请求:
# 使用 curl 测试
curl -X POST "http://localhost:8081/order/create" \
  -H "Content-Type: application/x-www-form-urlencoded" \
  -d "userId=1&productId=1&count=10&money=100"
正常情况测试:
  • 检查 seata_order 库的 t_order 表,应有一条 status=1 的订单记录
  • 检查 seata_storage 库的 t_storage 表,库存应减少10
  • 检查 seata_account 库的 t_account 表,余额应减少100
异常情况测试(验证回滚):
  1. 在账户服务的 decrease 方法中添加异常代码:
    throw new RuntimeException("模拟异常,测试回滚");
    
  2. 重新发送创建订单请求
  3. 验证结果:
    • 订单表中不应有新增记录或状态保持为0
    • 库存和账户余额应保持不变
    • 查看三个库的 undo_log 表,应有回滚日志记录
6. 核心原理说明
  1. 分布式事务标记:通过 @GlobalTransactional 注解标记分布式事务入口
  2. 事务协调
    • Seata 会生成全局事务ID (xid) 并在微服务间传递
    • 每个微服务的本地事务作为分支事务注册到 Seata Server
  3. 回滚机制
    • 正常情况:所有分支事务提交后,全局事务提交
    • 异常情况:Seata Server 协调所有分支事务回滚
    • 回滚依据:undo_log 表中记录的数据修改前后镜像

通过以上实现,三个微服务间的调用能够保证分布式事务的一致性,任何一个环节出现异常都会触发全局回滚。

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐