Saga分布式事务实战:Spring Boot微服务最终一致性方案

微服务架构中,跨服务的数据操作无法依赖本地事务保证ACID特性,分布式事务成为系统设计的核心挑战。Saga模式通过将长事务拆分为一系列本地事务,每个事务有对应的补偿操作,实现最终一致性。本文以Spring Boot为技术栈,演示Saga模式的完整实现方案。

Saga模式核心原理

Saga将一个分布式事务拆分为N个本地事务T1, T2, …, Tn,每个Ti有对应的补偿事务Ci。当Ti执行失败时,按逆序执行Ci-1, Ci-2, …, C1进行回滚。Saga有两种协调方式:

  • 编排式(Choreography):各服务通过事件订阅协作,无中心协调器
  • 编排式(Orchestration):由中央协调器统一调度各步骤执行和补偿

以电商订单流程为例,涉及订单服务、库存服务、支付服务三个微服务:

// 订单创建Saga流程
T1: 创建订单(订单服务)     C1: 取消订单
T2: 扣减库存(库存服务)     C2: 恢复库存
T3: 扣款支付(支付服务)     C3: 退款

// 正常流程: T1 → T2 → T3
// T3失败: 执行C2 → C1
// T2失败: 执行C1

编排式Saga实现

编排式Saga使用事件驱动,各服务发布事件并订阅其他服务的事件。以下基于Spring Boot + RabbitMQ实现:

// 事件定义
public class OrderEvents {
    public record OrderCreated(Long orderId, String userId, 
            List<OrderItem> items, BigDecimal amount) {}
    public record OrderCancelled(Long orderId, String reason) {}
}

public class InventoryEvents {
    public record InventoryDeducted(Long orderId, List<String> skus) {}
    public record InventoryDeductFailed(Long orderId, String reason) {}
}

public class PaymentEvents {
    public record PaymentCompleted(Long orderId, String transactionId) {}
    public record PaymentFailed(Long orderId, String reason) {}
}

// 订单服务 - 发布订单创建事件
@Service
public class OrderService {
    
    @Autowired
    private OrderRepository orderRepository;
    @Autowired
    private RabbitTemplate rabbitTemplate;

    @Transactional
    public Order createOrder(OrderRequest request) {
        Order order = Order.builder()
            .userId(request.getUserId())
            .items(request.getItems())
            .amount(request.getAmount())
            .status(OrderStatus.PENDING)
            .build();
        order = orderRepository.save(order);

        // 发布订单创建事件,触发Saga流程
        rabbitTemplate.convertAndSend("order.exchange", 
            "order.created", 
            new OrderEvents.OrderCreated(
                order.getId(), 
                order.getUserId(), 
                order.getItems(), 
                order.getAmount()
            ));
        
        return order;
    }

    // 补偿操作:取消订单
    @RabbitListener(queues = "order.cancel.queue")
    @Transactional
    public void handleOrderCancel(OrderEvents.OrderCancelled event) {
        Order order = orderRepository.findById(event.orderId())
            .orElseThrow();
        order.setStatus(OrderStatus.CANCELLED);
        order.setCancelReason(event.reason());
        orderRepository.save(order);
        log.info("订单 {} 已取消,原因: {}", event.orderId(), event.reason());
    }
}
// 库存服务 - 订阅订单创建事件,扣减库存
@Service
public class InventoryService {

    @Autowired
    private InventoryRepository inventoryRepository;
    @Autowired
    private RabbitTemplate rabbitTemplate;

    @RabbitListener(queues = "inventory.deduct.queue")
    @Transactional
    public void handleOrderCreated(OrderEvents.OrderCreated event) {
        try {
            for (OrderItem item : event.items()) {
                int updated = inventoryRepository.deductStock(
                    item.getSku(), item.getQuantity());
                if (updated == 0) {
                    throw new InsufficientStockException(
                        "SKU: " + item.getSku() + " 库存不足");
                }
            }
            // 发布库存扣减成功事件
            rabbitTemplate.convertAndSend("inventory.exchange",
                "inventory.deducted",
                new InventoryEvents.InventoryDeducted(
                    event.orderId(),
                    event.items().stream()
                        .map(OrderItem::getSku).toList()
                ));
        } catch (Exception e) {
            // 发布库存扣减失败事件,触发订单取消
            rabbitTemplate.convertAndSend("inventory.exchange",
                "inventory.deduct.failed",
                new InventoryEvents.InventoryDeductFailed(
                    event.orderId(), e.getMessage()
                ));
        }
    }

    // 补偿操作:恢复库存
    @RabbitListener(queues = "inventory.restore.queue")
    @Transactional
    public void handleInventoryRestore(PaymentEvents.PaymentFailed event) {
        // 恢复库存并发布取消订单事件
        inventoryRepository.restoreStock(event.orderId());
        rabbitTemplate.convertAndSend("order.exchange",
            "order.cancel",
            new OrderEvents.OrderCancelled(
                event.orderId(), "支付失败: " + event.reason()
            ));
    }
}

编排式Saga实现(Orchestration)

编排式Saga引入中央协调器(Orchestrator),由它统一管理事务状态和步骤调度。适合复杂流程,逻辑集中可控。使用Seata框架的Saga模式:

// Saga定义文件: order-saga.json
{
  "Name": "order-create-saga",
  "States": [
    {
      "Name": "CreateOrder",
      "Type": "ServiceTask",
      "ServiceName": "orderService",
      "ServiceMethod": "create",
      "CompensateState": "CancelOrder",
      "Input": [{"orderId": "$.{orderId}"}]
    },
    {
      "Name": "DeductInventory",
      "Type": "ServiceTask",
      "ServiceName": "inventoryService",
      "ServiceMethod": "deduct",
      "CompensateState": "RestoreInventory",
      "Input": [{"orderId": "$.{orderId}", "items": "$.{items}"}]
    },
    {
      "Name": "ProcessPayment",
      "Type": "ServiceTask",
      "ServiceName": "paymentService",
      "ServiceMethod": "pay",
      "CompensateState": "RefundPayment",
      "Input": [{"orderId": "$.{orderId}", "amount": "$.{amount}"}]
    },
    {
      "Name": "CancelOrder",
      "Type": "ServiceTask",
      "ServiceName": "orderService",
      "ServiceMethod": "cancel"
    },
    {
      "Name": "RestoreInventory",
      "Type": "ServiceTask",
      "ServiceName": "inventoryService",
      "ServiceMethod": "restore"
    },
    {
      "Name": "RefundPayment",
      "Type": "ServiceTask",
      "ServiceName": "paymentService",
      "ServiceMethod": "refund"
    }
  ]
}
// Spring Boot配置Seata Saga
@Configuration
public class SagaConfig {

    @Bean
    public StateMachineEngineImpl stateMachineEngine(
            SagaStateMachineEngine sagaEngine) {
        StateMachineEngineImpl engine = new StateMachineEngineImpl();
        engine.setSagaStateMachineEngine(sagaEngine);
        return engine;
    }
}

// Saga协调器服务
@Service
public class OrderSagaOrchestrator {

    @Autowired
    private StateMachineEngine stateMachineEngine;

    public void executeOrderSaga(OrderRequest request) {
        Map<String, Object> params = new HashMap<>();
        params.put("orderId", request.getOrderId());
        params.put("items", request.getItems());
        params.put("amount", request.getAmount());

        // 异步执行Saga
        stateMachineEngine.startWithBusinessKey(
            "order-create-saga",
            null,
            "ORDER_" + request.getOrderId(),
            params,
            (processId, state, result) -> {
                if (state == StateMachineEngine.State.FAILED) {
                    log.error("Saga执行失败,processId={}, orderId={}", 
                        processId, request.getOrderId());
                    // 触发告警通知
                    alertService.notifySagaFailure(
                        request.getOrderId(), result);
                } else if (state == StateMachineEngine.State.COMPLETED) {
                    log.info("Saga执行成功,orderId={}", 
                        request.getOrderId());
                }
            }
        );
    }
}

消息中间件可靠性保障

Saga模式依赖消息中间件传递事件,消息可靠性直接影响事务一致性。需要解决三个问题:消息丢失、消息重复、消息乱序。

// 本地消息表方案:确保消息可靠投递
@Service
public class ReliableMessagePublisher {

    @Autowired
    private MessageOutboxRepository outboxRepository;
    @Autowired
    private RabbitTemplate rabbitTemplate;

    // 与业务操作在同一本地事务中写入消息表
    @Transactional
    public void publishWithTransaction(String exchange, String routingKey, 
            Object payload, String businessId) {
        // 1. 业务操作(由调用方完成)
        
        // 2. 写入消息发件箱
        MessageOutbox message = MessageOutbox.builder()
            .businessId(businessId)
            .exchange(exchange)
            .routingKey(routingKey)
            .payload(JsonUtils.toJson(payload))
            .status(MessageStatus.PENDING)
            .retryCount(0)
            .build();
        outboxRepository.save(message);
    }

    // 定时任务扫描未发送消息
    @Scheduled(fixedDelay = 5000)
    public void scanAndPublish() {
        List<MessageOutbox> messages = outboxRepository
            .findPendingMessages(100);
        
        for (MessageOutbox msg : messages) {
            try {
                rabbitTemplate.convertAndSend(
                    msg.getExchange(),
                    msg.getRoutingKey(),
                    JsonUtils.fromJson(msg.getPayload(), Object.class)
                );
                msg.setStatus(MessageStatus.SENT);
                outboxRepository.save(msg);
            } catch (Exception e) {
                msg.incrementRetryCount();
                if (msg.getRetryCount() >= 5) {
                    msg.setStatus(MessageStatus.FAILED);
                    log.error("消息发送失败,businessId={}", msg.getBusinessId());
                }
                outboxRepository.save(msg);
            }
        }
    }
}

// 消费者幂等处理
@Component
public class IdempotentConsumer {

    @Autowired
    private RedisTemplate<String, String> redisTemplate;

    @RabbitListener(queues = "order.created.queue")
    public void handle(OrderEvents.OrderCreated event) {
        String key = "saga:consumed:" + event.orderId();
        
        // Redis SETNX保证幂等
        Boolean isFirst = redisTemplate.opsForValue()
            .setIfAbsent(key, "1", Duration.ofHours(24));
        
        if (Boolean.FALSE.equals(isFirst)) {
            log.warn("重复消息已忽略,orderId={}", event.orderId());
            return;
        }
        
        // 执行业务逻辑
        processOrder(event);
    }
}

服务治理与异常处理

Saga执行过程中可能出现的异常场景及处理策略:

// Saga状态持久化与恢复
@Service
public class SagaRecoveryService {

    @Autowired
    private SagaInstanceRepository sagaRepository;

    // 定时扫描中断的Saga实例
    @Scheduled(fixedDelay = 30000)
    public void recoverInterruptedSagas() {
        List<SagaInstance> staleSagas = sagaRepository
            .findStaleInstances(LocalDateTime.now().minusMinutes(5));
        
        for (SagaInstance saga : staleSagas) {
            try {
                switch (saga.getCurrentState()) {
                    case "DeductInventory" -> {
                        // 检查库存是否实际扣减
                        boolean deducted = inventoryService
                            .checkDeducted(saga.getOrderId());
                        if (deducted) {
                            // 继续执行下一步
                            saga.proceedTo("ProcessPayment");
                        } else {
                            // 执行补偿
                            saga.compensate();
                        }
                    }
                    case "ProcessPayment" -> {
                        boolean paid = paymentService
                            .checkPaid(saga.getOrderId());
                        if (paid) {
                            saga.complete();
                        } else {
                            saga.compensate();
                        }
                    }
                }
                sagaRepository.save(saga);
            } catch (Exception e) {
                log.error("Saga恢复失败,sagaId={}", saga.getId(), e);
                saga.incrementRecoveryAttempt();
                if (saga.getRecoveryAttempts() > 3) {
                    saga.setStatus(SagaStatus.MANUAL_INTERVENTION);
                    alertService.notifyManualIntervention(saga);
                }
                sagaRepository.save(saga);
            }
        }
    }
}

Saga模式选择建议:简单流程(2-3个步骤)使用编排式,开发成本低且无单点依赖;复杂流程(4个以上步骤、含分支逻辑)使用编排式,集中管理状态更易维护。无论哪种方式,消息中间件的可靠性保障和补偿操作的幂等设计都是最终一致性的关键保障。

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

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

相关推荐