Spring Boot微服务分布式事务Saga模式与Seata集成实战

微服务架构中的分布式事务是后端开发的核心挑战。单体架构中一个数据库事务即可保证ACID,拆分为多个服务后,跨库操作需要分布式事务协调。Seata作为成熟的分布式事务框架提供AT、TCC、Saga、XA四种模式。本文演示在Spring Boot微服务中集成Seata AT模式与Saga编排方案的完整实现。

分布式事务问题场景分析

以电商下单流程为例,涉及三个微服务:订单服务(创建订单)、库存服务(扣减库存)、账户服务(扣减余额)。三个服务各自操作独立数据库,需要保证三个操作要么全部成功要么全部回滚。

本地事务无法跨服务传播,网络超时或服务宕机会导致数据不一致。Seata通过TC(事务协调器)、TM(事务管理器)、RM(资源管理器)三个角色协调全局事务。AT模式自动生成回滚SQL,对业务代码侵入性最低。

Seata Server部署配置

下载部署Seata Server 2.0,使用Nacos作为注册中心和配置中心:

# application.yml (seata-server)
seata:
  config:
    type: nacos
    nacos:
      server-addr: 127.0.0.1:8848
      group: SEATA_GROUP
      namespace: seata
      data-id: seataServer.properties
  registry:
    type: nacos
    nacos:
      server-addr: 127.0.0.1:8848
      group: SEATA_GROUP
      namespace: seata
  store:
    mode: db
    db:
      datasource: druid
      db-type: mysql
      url: jdbc:mysql://127.0.0.1:3306/seata?useUnicode=true
      user: root
      password: ${DB_PASSWORD}

在Nacos中创建seataServer.properties配置,包含store模式、全局事务超时等参数。Seata Server使用数据库存储事务日志,确保Server重启后事务状态不丢失。生产环境建议部署集群模式(至少3节点)保障高可用。

微服务集成Seata AT模式

在每个微服务中添加Seata依赖:

<!-- pom.xml -->
<dependency>
  <groupId>com.alibaba.cloud</groupId>
  <artifactId>spring-cloud-starter-alibaba-seata</artifactId>
  <version>2023.0.1.0</version>
</dependency>
<dependency>
  <groupId>com.alibaba</groupId>
  <artifactId>seata-spring-boot-starter</artifactId>
  <version>2.0.0</version>
</dependency>

application.yml配置Seata客户端:

seata:
  enabled: true
  application-id: ${spring.application.name}
  tx-service-group: my_tx_group
  service:
    vgroup-mapping:
      my_tx_group: default
  registry:
    type: nacos
    nacos:
      server-addr: 127.0.0.1:8848
      group: SEATA_GROUP
      namespace: seata
  config:
    type: nacos
    nacos:
      server-addr: 127.0.0.1:8848
      group: SEATA_GROUP
      namespace: seata
  data-source-proxy-mode: AT

每个业务数据库需创建undo_log表,Seata AT模式通过此表记录数据变更前镜像用于自动回滚:

-- 在每个业务数据库执行
CREATE TABLE IF NOT EXISTS `undo_log` (
  `branch_id` BIGINT NOT NULL COMMENT 'branch transaction id',
  `xid` VARCHAR(128) NOT NULL COMMENT 'global transaction id',
  `context` VARCHAR(128) NOT NULL COMMENT 'undo_log context',
  `rollback_info` LONGBLOB NOT NULL COMMENT 'rollback info',
  `log_status` INT NOT NULL COMMENT '0:normal status,1:defense status',
  `log_created` DATETIME(6) NOT NULL COMMENT 'create datetime',
  `log_modified` DATETIME(6) NOT NULL COMMENT 'modify datetime',
  UNIQUE KEY `ux_undo_log` (`xid`, `branch_id`)
) ENGINE = InnoDB COMMENT = 'AT transaction mode undo table';

业务代码实现全局事务

订单服务作为TM(事务发起方),在Service层使用@GlobalTransactional注解开启全局事务:

// 订单服务 OrderService.java
@Service
@Slf4j
public class OrderService {
    
    @Resource
    private OrderMapper orderMapper;
    @Resource
    private StorageFeignClient storageClient;
    @Resource
    private AccountFeignClient accountClient;
    
    @GlobalTransactional(name = "createOrder", timeoutMills = 60000, rollbackFor = Exception.class)
    public Order createOrder(OrderDTO dto) {
        log.info("===== 开始创建订单, xid={}", RootContext.getXID());
        
        // 1. 创建订单(本地事务)
        Order order = new Order();
        order.setUserId(dto.getUserId());
        order.setProductId(dto.getProductId());
        order.setCount(dto.getCount());
        order.setMoney(dto.getMoney());
        order.setStatus(0);
        orderMapper.insert(order);
        log.info("订单创建成功, orderId={}", order.getId());
        
        // 2. 远程调用库存服务扣减库存
        Result storageResult = storageClient.decrease(dto.getProductId(), dto.getCount());
        if (!storageResult.isSuccess()) {
            throw new RuntimeException("库存扣减失败: " + storageResult.getMessage());
        }
        log.info("库存扣减成功");
        
        // 3. 远程调用账户服务扣减余额
        Result accountResult = accountClient.decrease(dto.getUserId(), dto.getMoney());
        if (!accountResult.isSuccess()) {
            throw new RuntimeException("余额扣减失败: " + accountResult.getMessage());
        }
        log.info("余额扣减成功");
        
        // 4. 更新订单状态
        order.setStatus(1);
        orderMapper.updateById(order);
        log.info("===== 订单创建完成, orderId={}", order.getId());
        
        return order;
    }
}

@GlobalTransactional标记的方法内所有数据库操作和远程调用都纳入同一全局事务。任一环节抛出异常,Seata TC会协调所有分支事务自动回滚。业务代码无需手动处理补偿逻辑。

Feign客户端传递全局事务ID

远程调用通过Feign传播xid,需添加Seata提供的请求拦截器:

// 库存服务 Feign接口
@FeignClient(name = "storage-service", contextId = "storage")
public interface StorageFeignClient {
    
    @PostMapping("/storage/decrease")
    Result decrease(@RequestParam("productId") Long productId, 
                          @RequestParam("count") Integer count);
}

// 账户服务 Feign接口
@FeignClient(name = "account-service", contextId = "account")
public interface AccountFeignClient {
    
    @PostMapping("/account/decrease")
    Result decrease(@RequestParam("userId") Long userId, 
                          @RequestParam("money") BigDecimal money);
}

Seata自动注册SeataFeignClient拦截器,在HTTP请求头中传递xid。分支服务通过Seata的DataSource代理自动加入到全局事务,生成undo_log记录。确保feign.hystrix.enabled=false或正确配置Hystrix上下文传递,避免线程切换导致xid丢失。

Saga模式处理长事务

AT模式适用于执行时间短的事务(秒级),对于执行时间较长的业务流程(分钟级或更长),Saga模式更合适。Saga通过编排正向操作和补偿操作实现最终一致性。

定义Saga状态机,使用Seata Saga状态机JSON配置:

// saga/create_order.json
{
  "Name": "createOrderSaga",
  "Comment": "创建订单Saga流程",
  "StartState": "CreateOrder",
  "States": {
    "CreateOrder": {
      "Type": "ServiceTask",
      "ServiceName": "orderService",
      "ServiceMethod": "create",
      "CompensateState": "CompensateOrder",
      "Next": "DeductStorage"
    },
    "DeductStorage": {
      "Type": "ServiceTask",
      "ServiceName": "storageService",
      "ServiceMethod": "decrease",
      "CompensateState": "CompensateStorage",
      "Next": "DeductAccount"
    },
    "DeductAccount": {
      "Type": "ServiceTask",
      "ServiceName": "accountService",
      "ServiceMethod": "decrease",
      "CompensateState": "CompensateAccount",
      "Next": "Succeed"
    },
    "CompensateOrder": {
      "Type": "ServiceTask",
      "ServiceName": "orderService",
      "ServiceMethod": "compensateCreate"
    },
    "CompensateStorage": {
      "Type": "ServiceTask",
      "ServiceName": "storageService",
      "ServiceMethod": "compensateDecrease"
    },
    "CompensateAccount": {
      "Type": "ServiceTask",
      "ServiceName": "accountService",
      "ServiceMethod": "compensateDecrease"
    },
    "Succeed": {
      "Type": "Succeed"
    }
  }
}

使用Saga状态机引擎执行流程:

// SagaOrderService.java
@Service
public class SagaOrderService {
    
    @Resource
    private StateMachineEngine stateMachineEngine;
    
    public String createOrderSaga(OrderDTO dto) {
        Map params = new HashMap<>();
        params.put("userId", dto.getUserId());
        params.put("productId", dto.getProductId());
        params.put("count", dto.getCount());
        params.put("money", dto.getMoney());
        
        StatusMachineInstance instance = stateMachineEngine.start(
            "createOrderSaga",
            null,
            params
        );
        
        if (instance.getStatus() == ExecutionStatus.SU) {
            return instance.getBusinessKey();
        } else {
            throw new RuntimeException("Saga事务执行失败: " + instance.getException());
        }
    }
}

Saga模式在任一步骤失败时,逆序执行已完成步骤的补偿操作。例如DeductAccount失败,会先执行CompensateStorage恢复库存,再执行CompensateOrder取消订单。补偿操作必须幂等,网络重试不会导致重复补偿。Saga模式放弃了ACID的隔离性,适合对一致性要求为最终一致的跨服务业务流程。

API接口规范与服务治理

分布式事务场景下的RestTemplate或Feign调用需配置合理的超时和重试策略,避免因网络抖动导致全局事务超时回滚:

// Feign超时配置
feign:
  client:
    config:
      default:
        connect-timeout: 3000
        read-timeout: 10000
  sentinel:
    enabled: true

# 重试配置
spring:
  cloud:
    loadbalancer:
      retry:
        enabled: true
        max-retries: 2
        retry-on-all-operations: false

全局事务超时时间需大于所有分支事务执行时间总和。Seata默认超时60秒,复杂业务流程需根据实际情况调大。超时后Seata TC主动发起全局回滚,未完成的分支事务会被强制回滚。消息中间件可用于异步解耦非核心流程,将其移出分布式事务范围,缩短事务路径。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot-wei-fu-wu-fen-bu-shi-shi-wu-saga-mo-shi-yu-seata/

(0)
小编小编
上一篇 12小时前
下一篇 12小时前

相关推荐