分布式事务Saga模式实战:Spring Boot微服务一致性保障方案

分布式事务问题背景与Saga模式原理

微服务架构中,一个业务操作跨越多个服务节点,每个服务维护独立数据库,传统本地ACID事务无法保证跨服务的数据一致性。分布式事务的核心挑战在于:部分服务执行成功而部分失败时,如何回滚已提交的本地事务。Saga模式将长事务拆分为一系列本地事务,每个本地事务有对应的补偿操作。正向执行时按顺序提交T1、T2、T3,任一步骤失败时反向执行C2、C1补偿已提交的事务。

Saga有两种协调方式:编排式(Orchestration)和协同式(Choreography)。编排式由中央协调器统一调度各服务执行顺序,适合复杂业务流程;协同式通过事件驱动各服务自主响应,适合简单流程。Spring Boot框架下,Seata中间件提供了开箱即用的Saga实现。

Seata Saga模式架构与部署

Seata是阿里巴巴开源的分布式事务解决方案,支持AT、TCC、SAGA、XA四种模式。Saga模式适用于长流程业务场景,如订单创建涉及库存扣减、优惠券核销、积分发放、支付处理等多个步骤。

Seata Server部署(Docker方式):

docker run -d --name seata-server \
  -p 8091:8091 \
  -p 7091:7091 \
  -e SEATA_IP=127.0.0.1 \
  -e STORE_MODE=db \
  -e STORE_DB_DRIVER_CLASS_NAME=com.mysql.cj.jdbc.Driver \
  -e STORE_DB_URL="jdbc:mysql://127.0.0.1:3306/seata?useUnicode=true&characterEncoding=utf8" \
  -e STORE_DB_USER=root \
  -e STORE_DB_PASSWORD='your_password' \
  seataio/seata-server:2.0.0

业务数据库需创建seata全局事务表和分支事务表:

CREATE TABLE IF NOT EXISTS `global_table` (
  `xid` varchar(128) NOT NULL,
  `transaction_id` bigint(20) DEFAULT NULL,
  `status` tinyint(4) NOT NULL,
  `application_id` varchar(32) DEFAULT NULL,
  `transaction_service_group` varchar(32) DEFAULT NULL,
  `transaction_name` varchar(128) DEFAULT NULL,
  `timeout` int(11) DEFAULT NULL,
  `begin_time` bigint(20) DEFAULT NULL,
  `application_data` varchar(2000) DEFAULT NULL,
  `gmt_create` datetime DEFAULT NULL,
  `gmt_modified` datetime DEFAULT NULL,
  PRIMARY KEY (`xid`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8;

CREATE TABLE IF NOT EXISTS `branch_table` (
  `branch_id` bigint(20) NOT NULL,
  `xid` varchar(128) NOT NULL,
  `transaction_id` bigint(20) DEFAULT NULL,
  `resource_group_id` varchar(32) DEFAULT NULL,
  `resource_id` varchar(256) DEFAULT NULL,
  `branch_type` varchar(8) DEFAULT NULL,
  `status` tinyint(4) DEFAULT NULL,
  `client_id` varchar(64) DEFAULT NULL,
  `application_data` varchar(2000) DEFAULT NULL,
  `gmt_create` datetime DEFAULT NULL,
  `gmt_modified` datetime DEFAULT NULL,
  PRIMARY KEY (`branch_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8;

Spring Boot集成Seata Saga实战

以电商订单流程为例,定义Saga编排流程。订单服务作为入口,依次调用库存服务、优惠券服务、积分服务。

Maven依赖配置:

<dependency>
  <groupId>io.seata</groupId>
  <artifactId>seata-spring-boot-starter</artifactId>
  <version>2.0.0</version>
</dependency>
<dependency>
  <groupId>io.seata</groupId>
  <artifactId>seata-saga-engine</artifactId>
  <version>2.0.0</version>
</dependency>

application.yml配置:

seata:
  enabled: true
  application-id: order-service
  tx-service-group: my_tx_group
  service:
    vgroup-mapping:
      my_tx_group: default
    grouplist:
      default: 127.0.0.1:8091
  saga:
    enabled: true
    state-machine:
      enable-async: false
      trans-operation-timeout: 1800000

Saga状态机JSON定义(定义正向和补偿流程):

{
  "Name": "createOrderSaga",
  "Comment": "订单创建Saga流程",
  "StartState": "deductInventory",
  "States": {
    "deductInventory": {
      "Type": "ServiceTask",
      "ServiceName": "inventoryService",
      "ServiceMethod": "deduct",
      "CompensateState": "compensateInventory",
      "Next": "useCoupon"
    },
    "compensateInventory": {
      "Type": "ServiceTask",
      "ServiceName": "inventoryService",
      "ServiceMethod": "addBack"
    },
    "useCoupon": {
      "Type": "ServiceTask",
      "ServiceName": "couponService",
      "ServiceMethod": "use",
      "CompensateState": "compensateCoupon",
      "Next": "addPoints"
    },
    "compensateCoupon": {
      "Type": "ServiceTask",
      "ServiceName": "couponService",
      "ServiceMethod": "restore"
    },
    "addPoints": {
      "Type": "ServiceTask",
      "ServiceName": "pointsService",
      "ServiceMethod": "add",
      "CompensateState": "compensatePoints",
      "Next": "succeed"
    },
    "compensatePoints": {
      "Type": "ServiceTask",
      "ServiceName": "pointsService",
      "ServiceMethod": "deduct"
    },
    "succeed": {
      "Type": "Succeed"
    }
  }
}

Java服务实现与触发:

@Service
public class OrderSagaService {

    @Autowired
    private StateMachineEngine stateMachineEngine;

    public String createOrder(OrderDTO order) {
        Map<String, Object> params = new HashMap<>();
        params.put("orderId", order.getOrderId());
        params.put("productId", order.getProductId());
        params.put("quantity", order.getQuantity());
        params.put("couponId", order.getCouponId());
        params.put("userId", order.getUserId());

        StateMachineInstance instance = stateMachineEngine.start(
            "createOrderSaga",
            BusinessType.ORDER.name(),
            params
        );

        if (instance.getStatus() == ExecutionStatus.SU) {
            return "订单创建成功: " + order.getOrderId();
        } else {
            throw new RuntimeException("订单创建失败,已执行补偿: " 
                + instance.getException());
        }
    }
}

@Service
public class InventoryServiceImpl implements InventoryService {

    @Autowired
    private InventoryMapper inventoryMapper;

    @Override
    public boolean deduct(Map<String, Object> params) {
        String productId = (String) params.get("productId");
        Integer quantity = (Integer) params.get("quantity");
        int rows = inventoryMapper.deductStock(productId, quantity);
        if (rows == 0) {
            throw new RuntimeException("库存不足: " + productId);
        }
        return true;
    }

    @Override
    public boolean addBack(Map<String, Object> params) {
        String productId = (String) params.get("productId");
        Integer quantity = (Integer) params.get("quantity");
        inventoryMapper.addBackStock(productId, quantity);
        return true;
    }
}

消息中间件在Saga中的可靠性保障

Saga模式中补偿操作的触发依赖协调器与服务间的通信可靠性。引入消息中间件(如RocketMQ)作为服务间通信通道,确保补偿指令不丢失。RocketMQ事务消息机制保证本地事务与消息发送的原子性:

@Service
public class OrderTransactionListener implements TransactionListener {

    @Autowired
    private OrderService orderService;

    @Override
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        try {
            OrderDTO order = JSON.parseObject(
                new String(msg.getBody()), OrderDTO.class);
            orderService.createOrder(order);
            return LocalTransactionState.COMMIT_MESSAGE;
        } catch (Exception e) {
            return LocalTransactionState.ROLLBACK_MESSAGE;
        }
    }

    @Override
    public LocalTransactionState checkLocalTransaction(MessageExt msg) {
        String orderId = msg.getKeys();
        Order order = orderService.getById(orderId);
        if (order != null && order.getStatus() == OrderStatus.CREATED) {
            return LocalTransactionState.COMMIT_MESSAGE;
        }
        return LocalTransactionState.UNKNOW;
    }
}

服务治理与幂等性设计

Saga补偿操作必须保证幂等性,因为网络超时可能导致重试。补偿方法通过状态机版本号判断是否已执行,避免重复补偿导致数据不一致。业务中台建设中,建议为每个Saga步骤维护状态表,记录正向和补偿操作的执行状态:

CREATE TABLE saga_step_log (
  id BIGINT AUTO_INCREMENT PRIMARY KEY,
  saga_id VARCHAR(64) NOT NULL,
  step_name VARCHAR(64) NOT NULL,
  step_type VARCHAR(16) NOT NULL COMMENT 'FORWARD/COMPENSATION',
  status VARCHAR(16) NOT NULL COMMENT 'INIT/RUNNING/SUCCESS/FAIL',
  retry_count INT DEFAULT 0,
  request_data TEXT,
  response_data TEXT,
  created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
  UNIQUE KEY uk_saga_step (saga_id, step_name, step_type)
);

补偿方法执行前检查该步骤是否已成功补偿,已成功则直接返回,避免重复执行。API接口规范方面,Saga协调器调用各服务接口统一使用POST方法,请求体携带sagaId和stepName字段,服务端通过这两个字段实现幂等控制。

高并发设计场景下,Saga模式需配合分布式锁防止同一业务流程并发执行。使用Redisson分布式锁在订单维度加锁,确保同一订单不会同时触发多个Saga实例。锁粒度尽量细化到业务实体级别(如orderId),避免全局锁影响吞吐量。

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

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

相关推荐