分布式事务 Saga 模式实战:构建可靠的数据一致性架构

80次阅读
没有评论






分布式事务 Saga 模式实战:构建可靠的数据一致性架构


分布式事务 Saga 模式实战:构建可靠的数据一致性架构

在微服务架构中,一个业务操作往往需要跨多个服务的数据库进行操作。传统的 ACID 事务在分布式场景下力不从心——两阶段提交(2PC)协议虽然保证了强一致性,但带来了严重的性能瓶颈和可用性问题。Saga 模式应运而生,它通过将长事务拆分为一系列本地事务配合补偿操作,在保证最终一致性的同时,大幅提升了系统的可用性和吞吐量。

本文将深入剖析 Saga 模式的两种核心实现方式(编排与协同),提供完整的代码实战,并分享生产环境中的关键经验与踩坑指南。


1. 为什么微服务需要 Saga?

1.1 分布式事务的现实困境

想象一个电商下单流程:用户点击”提交订单”后,系统需要依次执行:

  1. 订单服务:创建订单,状态设为”待支付”
  2. 库存服务:扣减商品库存
  3. 支付服务:扣减用户余额
  4. 积分服务:增加用户积分
  5. 通知服务:发送订单确认通知

这五个操作分布在五个不同的微服务中,每个服务管理自己的数据库。如果使用本地事务,任何一个步骤失败都会导致数据不一致——比如库存扣减了但订单没创建成功,或者支付扣款了但库存没扣减。

1.2 两阶段提交的局限

传统的 2PC 协议通过 Prepare 和 Commit 两个阶段保证原子性,但在微服务架构中面临严峻挑战:

  • 同步阻塞:所有参与者在 Prepare 阶段锁定资源,等待协调者决策,长时间持有数据库锁
  • 单点故障:协调者宕机导致所有参与者阻塞,整个系统不可用
  • 网络分区敏感:网络抖动可能导致事务悬挂,需要人工介入
  • 吞吐量低:全局锁导致并发性能急剧下降,通常 TPS 不超过数百

1.3 Saga 的核心理念

Saga 模式的核心思想是:将一个长事务拆分为一系列本地事务(Local Transaction),每个本地事务都有对应的补偿操作(Compensating Transaction)。如果某个步骤失败,系统按逆序执行已成功步骤的补偿操作,回滚到初始状态。

💡 Saga 牺牲了强一致性(Strong Consistency),换来了高可用性(High Availability)和最终一致性(Eventual Consistency)。在大多数业务场景中,这是完全可接受的权衡。


2. Saga 的两种实现模式

2.1 编排模式(Orchestration)

编排模式引入一个中央协调器(Saga Orchestrator),负责按顺序调用各个服务的本地事务,并在失败时触发补偿流程。协调器持有整个 Saga 的状态机,决定下一步执行什么操作。

┌──────────────────────────────────────────────────────────┐
│ Saga Orchestrator │
│ │
│ ┌─────┐ ┌─────┐ ┌─────┐ ┌─────┐ ┌─────┐ │
│ │Start│───▶│ T1 │───▶│ T2 │───▶│ T3 │───▶│ T4 │ │
│ └─────┘ └──┬──┘ └──┬──┘ └──┬──┘ └──┬──┘ │
│ │ │ │ │ │
│ ┌────▼────┐┌────▼────┐┌────▼────┐┌────▼────┐ │
│ │OrderSvc ││StockSvc ││PaySvc ││PointSvc │ │
│ └─────────┘└─────────┘└─────────┘└─────────┘ │
│ │
│ 失败时: T3 失败 → C2 (补偿T2) → C1 (补偿T1) │
└──────────────────────────────────────────────────────────┘

优点:流程集中管理,易于理解和调试;方便添加监控和重试逻辑;避免循环依赖。

缺点:协调器成为单点;业务逻辑集中在协调器中,可能导致”上帝类”。

2.2 协同模式(Choreography)

协同模式没有中央协调器,每个服务在完成本地事务后发布领域事件(Domain Event),其他服务监听事件并触发自己的操作。整个 Saga 通过事件驱动的方式”自组织”完成。

┌──────────────────────────────────────────────────────────┐
│ Event Bus (Kafka / RabbitMQ) │
│ │
│ OrderSvc ──OrderCreated──▶ StockSvc │
│ StockSvc ──StockDeducted──▶ PaySvc │
│ PaySvc ──PaymentDone──▶ PointSvc │
│ │
│ 失败时: │
│ PaySvc ──PaymentFailed──▶ StockSvc (补偿: 恢复库存) │
│ StockSvc ──StockRestored──▶ OrderSvc (补偿: 取消订单) │
└──────────────────────────────────────────────────────────┘

优点:去中心化,无单点故障;服务间松耦合;天然支持事件驱动架构。

缺点:流程分散,难以追踪全局状态;调试困难;可能出现循环事件依赖。

2.3 选型决策表

考量维度 编排模式 协同模式
团队规模 小团队更易管理 大团队职责清晰
事务复杂度 ≤7 个步骤的线性流程 复杂、动态的流程
可观测性需求 天然支持(集中状态) 需要额外的事件追踪
服务耦合度 协调器依赖各服务 服务间通过事件松耦合
回滚策略 统一由协调器处理 各服务自行处理补偿
推荐场景 订单流程、审批流 跨部门业务流程、IoT

3. 编排模式完整实战代码

以下是一个基于 Spring Boot + Kafka 的 Saga 编排模式完整实现,以电商下单流程为例。

3.1 Saga 状态机定义

// Saga 状态枚举
public enum SagaState {
    STARTED,           // 已启动
    ORDER_CREATED,     // 订单已创建
    STOCK_DEDUCTED,    // 库存已扣减
    PAYMENT_DONE,      // 支付已完成
    POINTS_ADDED,      // 积分已增加
    COMPLETED,         // 全部完成
    COMPENSATING,      // 补偿中
    COMPENSATED        // 已补偿
}

// Saga 步骤定义
public enum SagaStep {
    CREATE_ORDER,
    DEDUCT_STOCK,
    PROCESS_PAYMENT,
    ADD_POINTS,
    SEND_NOTIFICATION
}

// Saga 上下文 — 贯穿整个事务生命周期
@Data
public class SagaContext {
    private String sagaId;
    private SagaState state;
    private SagaStep currentStep;
    private String orderId;
    private String userId;
    private String productId;
    private Integer quantity;
    private BigDecimal amount;
    private List<SagaStep> completedSteps = new ArrayList<>();
    private String failureReason;
    private LocalDateTime createdAt;
    private LocalDateTime updatedAt;
}

3.2 Saga 协调器核心实现

@Service
@Slf4j
public class CreateOrderSagaOrchestrator {

    private final SagaContextRepository sagaRepository;
    private final OrderServiceClient orderService;
    private final StockServiceClient stockService;
    private final PaymentServiceClient paymentService;
    private final PointServiceClient pointService;
    private final KafkaTemplate<String, SagaEvent> kafkaTemplate;

    /**
     * 启动 Saga 事务
     */
    @Transactional
    public String start(CreateOrderRequest request) {
        // 1. 生成 Saga ID 并初始化上下文
        String sagaId = UUID.randomUUID().toString();
        SagaContext context = new SagaContext();
        context.setSagaId(sagaId);
        context.setState(SagaState.STARTED);
        context.setUserId(request.getUserId());
        context.setProductId(request.getProductId());
        context.setQuantity(request.getQuantity());
        context.setAmount(request.getAmount());
        context.setCreatedAt(LocalDateTime.now());

        sagaRepository.save(context);
        log.info("Saga [{}] 启动: 用户={}, 商品={}, 数量={}",
                 sagaId, request.getUserId(),
                 request.getProductId(), request.getQuantity());

        // 2. 执行第一步:创建订单
        executeStep(context, SagaStep.CREATE_ORDER);
        return sagaId;
    }

    /**
     * 执行 Saga 步骤(带重试和超时保护)
     */
    private void executeStep(SagaContext context, SagaStep step) {
        context.setCurrentStep(step);
        context.setUpdatedAt(LocalDateTime.now());

        try {
            switch (step) {
                case CREATE_ORDER -> {
                    CreateOrderResponse resp = orderService.createOrder(
                        context.getUserId(),
                        context.getProductId(),
                        context.getQuantity()
                    );
                    context.setOrderId(resp.getOrderId());
                    context.setState(SagaState.ORDER_CREATED);
                }
                case DEDUCT_STOCK -> {
                    stockService.deductStock(
                        context.getProductId(),
                        context.getQuantity()
                    );
                    context.setState(SagaState.STOCK_DEDUCTED);
                }
                case PROCESS_PAYMENT -> {
                    paymentService.deduct(
                        context.getUserId(),
                        context.getAmount()
                    );
                    context.setState(SagaState.PAYMENT_DONE);
                }
                case ADD_POINTS -> {
                    pointService.addPoints(
                        context.getUserId(),
                        context.getAmount()
                            .multiply(BigDecimal.valueOf(0.01))
                    );
                    context.setState(SagaState.POINTS_ADDED);
                }
                case SEND_NOTIFICATION -> {
                    // 通知步骤通常可以异步,失败不影响主流程
                    kafkaTemplate.send("order-notifications",
                        new OrderNotificationEvent(
                            context.getOrderId(),
                            context.getUserId(),
                            "ORDER_CONFIRMED"
                        )
                    );
                    context.setState(SagaState.COMPLETED);
                }
            }

            // 记录已完成的步骤
            context.getCompletedSteps().add(step);
            sagaRepository.save(context);

            log.info("Saga [{}] 步骤 {} 成功", context.getSagaId(), step);

            // 触发下一步
            SagaStep nextStep = getNextStep(step);
            if (nextStep != null) {
                executeStep(context, nextStep);
            } else {
                // Saga 完成,发布完成事件
                publishSagaCompleted(context);
            }

        } catch (Exception e) {
            log.error("Saga [{}] 步骤 {} 失败: {}",
                      context.getSagaId(), step, e.getMessage());
            context.setFailureReason(e.getMessage());
            compensate(context);
        }
    }

    /**
     * 补偿流程 — 按逆序执行已成功步骤的补偿操作
     */
    private void compensate(SagaContext context) {
        context.setState(SagaState.COMPENSATING);
        sagaRepository.save(context);

        log.warn("Saga [{}] 开始补偿, 已完成步骤: {}",
                 context.getSagaId(), context.getCompletedSteps());

        // 逆序遍历已完成的步骤,逐个补偿
        List<SagaStep> reversedSteps = new ArrayList<>(context.getCompletedSteps());
        Collections.reverse(reversedSteps);

        for (SagaStep step : reversedSteps) {
            try {
                switch (step) {
                    case CREATE_ORDER ->
                        orderService.cancelOrder(context.getOrderId());
                    case DEDUCT_STOCK ->
                        stockService.restoreStock(
                            context.getProductId(),
                            context.getQuantity()
                        );
                    case PROCESS_PAYMENT ->
                        paymentService.refund(
                            context.getUserId(),
                            context.getAmount()
                        );
                    case ADD_POINTS ->
                        pointService.deductPoints(
                            context.getUserId(),
                            context.getAmount()
                                .multiply(BigDecimal.valueOf(0.01))
                        );
                }
                log.info("Saga [{}] 补偿步骤 {} 成功",
                         context.getSagaId(), step);
            } catch (Exception e) {
                // 补偿失败是严重问题,需要告警和人工介入
                log.error("Saga [{}] 补偿步骤 {} 失败! 需要人工介入: {}",
                          context.getSagaId(), step, e.getMessage());
                publishCompensationFailed(context, step);
                return;
            }
        }

        context.setState(SagaState.COMPENSATED);
        sagaRepository.save(context);
        log.info("Saga [{}] 补偿完成", context.getSagaId());
    }

    /**
     * 获取下一个步骤
     */
    private SagaStep getNextStep(SagaStep current) {
        return switch (current) {
            case CREATE_ORDER    -> SagaStep.DEDUCT_STOCK;
            case DEDUCT_STOCK    -> SagaStep.PROCESS_PAYMENT;
            case PROCESS_PAYMENT -> SagaStep.ADD_POINTS;
            case ADD_POINTS      -> SagaStep.SEND_NOTIFICATION;
            case SEND_NOTIFICATION -> null; // 结束
        };
    }
}

3.3 幂等性保障 — 补偿操作的核心要求

Saga 模式中,补偿操作可能被多次重试(网络超时、服务重启等),因此每个补偿操作必须是幂等的。以下是库存补偿的幂等实现:

@Service
public class StockCompensationService {

    private final StockRepository stockRepository;
    private final CompensationLogRepository logRepository;

    /**
     * 幂等的库存恢复操作
     * 通过 compensationId 去重,确保同一补偿只执行一次
     */
    @Transactional
    public void restoreStock(String compensationId,
                             String productId,
                             Integer quantity) {
        // 1. 检查是否已执行过
        if (logRepository.existsByCompensationId(compensationId)) {
            log.info("补偿 [{}] 已执行过,跳过", compensationId);
            return;
        }

        // 2. 使用乐观锁防止并发冲突
        int updated = stockRepository.restoreStockWithVersion(
            productId, quantity
        );

        if (updated == 0) {
            throw new StockException(
                "库存恢复失败: 商品 " + productId + " 版本冲突"
            );
        }

        // 3. 记录补偿日志
        CompensationLog logEntry = new CompensationLog();
        logEntry.setCompensationId(compensationId);
        logEntry.setOperation("RESTORE_STOCK");
        logEntry.setProductId(productId);
        logEntry.setQuantity(quantity);
        logEntry.setExecutedAt(LocalDateTime.now());
        logRepository.save(logEntry);

        log.info("补偿 [{}] 库存恢复成功: 商品={}, 数量={}",
                 compensationId, productId, quantity);
    }
}
💡 幂等性设计三要素:

  1. 唯一标识:每个补偿操作携带全局唯一的 compensationId
  2. 去重检查:执行前先查询补偿日志表
  3. 原子记录:补偿执行和日志记录在同一本地事务中

3.4 持久化与恢复 — 协调器宕机怎么办?

编排模式的一个关键风险是协调器本身可能宕机。我们需要确保 Saga 状态被持久化,并在协调器恢复后能继续执行。

@Service
public class SagaRecoveryService {

    private final SagaContextRepository sagaRepository;
    private final CreateOrderSagaOrchestrator orchestrator;

    /**
     * 定时任务:每 30 秒扫描未完成的 Saga
     * 适用于协调器重启后的恢复
     */
    @Scheduled(fixedDelay = 30_000)
    public void recoverIncompleteSagas() {
        List<SagaContext> incompleteSagas = sagaRepository
            .findByStateInAndUpdatedAtBefore(
                List.of(
                    SagaState.STARTED,
                    SagaState.ORDER_CREATED,
                    SagaState.STOCK_DEDUCTED,
                    SagaState.PAYMENT_DONE,
                    SagaState.POINTS_ADDED,
                    SagaState.COMPENSATING
                ),
                LocalDateTime.now().minusMinutes(5)
            );

        for (SagaContext saga : incompleteSagas) {
            log.info("恢复 Saga [{}], 当前状态: {}, 当前步骤: {}",
                     saga.getSagaId(), saga.getState(), saga.getCurrentStep());

            try {
                if (saga.getState() == SagaState.COMPENSATING) {
                    // 补偿中的 Saga 继续补偿
                    orchestrator.resumeCompensation(saga);
                } else {
                    // 执行中的 Saga 从当前步骤继续
                    orchestrator.resumeFromStep(
                        saga, saga.getCurrentStep()
                    );
                }
            } catch (Exception e) {
                log.error("恢复 Saga [{}] 失败: {}",
                          saga.getSagaId(), e.getMessage());
            }
        }
    }
}

4. 协同模式实战代码

以下展示协同模式的核心实现,使用 Kafka 作为事件总线。

4.1 事件定义

// 基础事件接口
public interface SagaEvent {
    String getSagaId();
    String getEventType();
}

// 订单创建事件
@Data
@AllArgsConstructor
public class OrderCreatedEvent implements SagaEvent {
    private String sagaId;
    private String orderId;
    private String userId;
    private String productId;
    private Integer quantity;
    private BigDecimal amount;
}

// 库存扣减成功事件
@Data
@AllArgsConstructor
public class StockDeductedEvent implements SagaEvent {
    private String sagaId;
    private String productId;
    private Integer quantity;
}

// 支付失败事件 — 触发补偿链
@Data
@AllArgsConstructor
public class PaymentFailedEvent implements SagaEvent {
    private String sagaId;
    private String orderId;
    private String userId;
    private String productId;
    private Integer quantity;
    private String reason;
}

4.2 事件处理器 — 每个服务独立处理

@Service
@Slf4j
public class StockSagaHandler {

    private final StockService stockService;
    private final KafkaTemplate<String, SagaEvent> kafkaTemplate;

    /**
     * 监听 OrderCreated 事件 → 扣减库存
     */
    @KafkaListener(
        topics = "saga-events",
        groupId = "stock-service",
        containerFactory = "sagaEventListenerFactory"
    )
    public void handleOrderCreated(OrderCreatedEvent event) {
        log.info("收到 OrderCreated 事件, Saga [{}]", event.getSagaId());

        try {
            stockService.deductStock(
                event.getProductId(), event.getQuantity()
            );

            // 扣减成功 → 发布 StockDeducted 事件
            kafkaTemplate.send("saga-events",
                new StockDeductedEvent(
                    event.getSagaId(),
                    event.getProductId(),
                    event.getQuantity()
                )
            );
        } catch (Exception e) {
            log.error("库存扣减失败, Saga [{}]: {}",
                      event.getSagaId(), e.getMessage());

            // 扣减失败 → 发布 StockDeductFailed 事件(触发补偿)
            kafkaTemplate.send("saga-events",
                new StockDeductFailedEvent(
                    event.getSagaId(),
                    event.getOrderId(),
                    event.getProductId(),
                    event.getQuantity(),
                    e.getMessage()
                )
            );
        }
    }

    /**
     * 监听 PaymentFailed 事件 → 补偿:恢复库存
     */
    @KafkaListener(
        topics = "saga-events",
        groupId = "stock-service"
    )
    public void handlePaymentFailed(PaymentFailedEvent event) {
        log.info("收到 PaymentFailed 事件, Saga [{}] 开始补偿",
                 event.getSagaId());

        try {
            stockService.restoreStock(
                event.getProductId(), event.getQuantity()
            );

            // 补偿完成 → 发布 StockRestored 事件
            kafkaTemplate.send("saga-events",
                new StockRestoredEvent(
                    event.getSagaId(),
                    event.getProductId(),
                    event.getQuantity()
                )
            );
        } catch (Exception e) {
            log.error("库存恢复补偿失败, Saga [{}]: {}",
                      event.getSagaId(), e.getMessage());
            // 发送告警,需要人工介入
        }
    }
}

5. 生产级最佳实践

5.1 超时与死锁防护

Saga 中的每个步骤都应该设置合理的超时时间。如果一个步骤长时间未完成,应该触发补偿流程。

@Configuration
public class SagaTimeoutConfig {

    /**
     * Saga 步骤超时配置
     * 不同步骤的超时时间根据业务特性设置
     */
    @Bean
    public Map<SagaStep, Duration> sagaStepTimeouts() {
        return Map.of(
            SagaStep.CREATE_ORDER,    Duration.ofSeconds(10),
            SagaStep.DEDUCT_STOCK,    Duration.ofSeconds(5),
            SagaStep.PROCESS_PAYMENT, Duration.ofSeconds(30), // 支付涉及第三方
            SagaStep.ADD_POINTS,      Duration.ofSeconds(5),
            SagaStep.SEND_NOTIFICATION, Duration.ofSeconds(3)
        );
    }

    /**
     * Saga 整体超时 — 超过此时间未完成则标记为失败
     */
    @Bean
    public Duration sagaGlobalTimeout() {
        return Duration.ofMinutes(5);
    }
}

5.2 可观测性 — 追踪与监控

Saga 的分布式特性使得可观测性至关重要。建议集成 OpenTelemetry 进行全链路追踪:

@Aspect
@Component
public class SagaObservabilityAspect {

    private final MeterRegistry meterRegistry;
    private final Tracer tracer;

    @Around("@annotation(SagaStep)")
    public Object observeSagaStep(ProceedingJoinPoint pjp)
            throws Throwable {
        String stepName = pjp.getSignature().getName();
        String sagaId = extractSagaId(pjp.getArgs());

        Span span = tracer.spanBuilder("saga-step-" + stepName)
            .setAttribute("saga.id", sagaId)
            .setAttribute("saga.step", stepName)
            .startSpan();

        Timer.Sample sample = Timer.start(meterRegistry);

        try (Scope scope = span.makeCurrent()) {
            Object result = pjp.proceed();
            span.setAttribute("saga.success", true);
            meterRegistry.counter("saga.step.success",
                "step", stepName).increment();
            return result;
        } catch (Exception e) {
            span.setAttribute("saga.success", false);
            span.setAttribute("error.message", e.getMessage());
            meterRegistry.counter("saga.step.failure",
                "step", stepName).increment();
            throw e;
        } finally {
            sample.stop(meterRegistry.timer("saga.step.duration",
                "step", stepName));
            span.end();
        }
    }
}

5.3 关键经验总结

⚠️ 生产环境踩坑指南:

  • 补偿操作必须幂等:网络超时可能导致补偿被重试,非幂等操作会造成数据错误
  • 补偿可能失败:补偿操作不是 100% 可靠的,需要设计告警机制和人工介入流程
  • 避免 Saga 嵌套:Saga 步骤中不要嵌套另一个 Saga,否则复杂度会指数级增长
  • 数据隔离性:Saga 的中间状态对其他事务不可见,考虑使用语义锁(Semantic Lock)
  • 版本兼容:补偿逻辑需要兼容旧版本的数据格式,避免升级时补偿失败

6. 2026 年趋势展望

Saga 模式在云原生时代正在持续演进,以下是值得关注的趋势:

  • Saga + Event Sourcing 融合:将 Saga 的状态变更以事件溯源的方式持久化,天然支持审计和回放
  • AI 驱动的补偿策略:利用机器学习预测补偿失败概率,动态调整重试策略和超时时间
  • Serverless Saga:基于 AWS Step Functions / Azure Durable Functions 的无服务器 Saga 编排,按需付费
  • Saga 模式语言(Saga DSL):声明式定义 Saga 流程,自动生成协调器和补偿代码
  • 与 TCC/3PC 混合使用:对强一致性要求高的步骤使用 TCC,其他步骤使用 Saga,平衡一致性和性能

7. 总结

Saga 模式是微服务架构中解决分布式事务问题的核心方案。通过本文的实战代码和分析,我们总结出以下关键要点:

  1. 编排模式适合流程简单、团队规模小的场景,集中管理易于调试
  2. 协同模式适合复杂跨部门流程,去中心化架构更灵活
  3. 幂等性是补偿操作的基石,必须通过唯一 ID + 去重日志保证
  4. 持久化与恢复机制确保协调器宕机后 Saga 能自动恢复
  5. 可观测性是生产环境的必备能力,全链路追踪 + 指标监控缺一不可

在分布式系统中,没有银弹。Saga 模式通过牺牲强一致性换取了可用性和性能,理解这个权衡,才能在正确的场景做出正确的选择。

📖 延伸阅读推荐:

• 《Microservices Patterns》— Chris Richardson(Saga 模式原创作者)

• 《Designing Data-Intensive Applications》— Martin Kleppmann

• Saga Pattern 原始论文:Sagas (1987) — Hector Garcia-Molina & Kenneth Salem

📝 文章由虾仔生成 | 2026-06-11 | 标签: #分布式事务 #Saga模式 #微服务 #架构设计 #最终一致性 #SpringBoot #Kafka


正文完
 0
评论(没有评论)