分布式事务Saga模式在Spring Boot微服务中的落地实现

分布式事务的核心矛盾与Saga的解法

微服务拆分后,一个业务操作往往跨越多个服务,每个服务拥有独立数据库。传统的两阶段提交(2PC)在微服务场景下因锁持有时间长、可用性差而不再适用。Saga模式的核心思路是:将长事务拆成多个本地事务,每个本地事务提交后发送事件触发下一步;任意步骤失败时,执行对应步骤的补偿操作回滚已完成的部分。

Saga编排模式:集中控制 vs 事件驱动

Saga有两种实现风格:编排式(Orchestration)由一个中心协调器控制流程,事件驱动式(Choreography)由各服务监听事件自主决定下一步。

编排式适合流程明确的业务(如订单创建),事件驱动式适合解耦度高的场景(如通知触达)。生产环境推荐编排式,因为流程可视、调试方便。

// Spring Boot编排式Saga实现
// 1. 定义Saga步骤
public class OrderSagaDefinition {

    private final List<SagaStep> steps;

    public OrderSagaDefinition() {
        this.steps = List.of(
            SagaStep.builder()
                .name("createOrder")
                .action("order-service", "POST /orders")
                .compensateAction("order-service", "DELETE /orders/{id}")
                .build(),
            SagaStep.builder()
                .name("reserveStock")
                .action("inventory-service", "POST /stock/reserve")
                .compensateAction("inventory-service", "POST /stock/release")
                .build(),
            SagaStep.builder()
                .name("processPayment")
                .action("payment-service", "POST /payments")
                .compensateAction("payment-service", "POST /payments/{id}/refund")
                .build(),
            SagaStep.builder()
                .name("confirmOrder")
                .action("order-service", "PUT /orders/{id}/confirm")
                .compensateAction(null, null)  // 最后一步无需补偿
                .build()
        );
    }
}

Saga协调器的核心实现

// SagaOrchestrator.java
@Service
@Slf4j
public class SagaOrchestrator {

    private final RestTemplate restTemplate;
    private final SagaInstanceRepository sagaRepo;

    public SagaOrchestrator(RestTemplate restTemplate, SagaInstanceRepository sagaRepo) {
        this.restTemplate = restTemplate;
        this.sagaRepo = sagaRepo;
    }

    public SagaResult execute(SagaDefinition definition, Map<String, Object> context) {
        String sagaId = UUID.randomUUID().toString();
        SagaInstance saga = new SagaInstance(sagaId, "RUNNING", System.currentTimeMillis());
        sagaRepo.save(saga);

        List<SagaStep> steps = definition.getSteps();
        int completedStep = -1;

        // 正向执行
        for (int i = 0; i < steps.size(); i++) {
            SagaStep step = steps.get(i);
            try {
                log.info("Saga[{}] executing step: {}", sagaId, step.getName());
                Object result = callService(step.getActionEndpoint(), step.getActionMethod(), context);
                context.put(step.getName() + "_result", result);
                completedStep = i;
                saga.setCurrentStep(i);
                sagaRepo.save(saga);
            } catch (Exception e) {
                log.error("Saga[{}] step {} failed: {}", sagaId, step.getName(), e.getMessage());
                // 触发补偿
                compensate(steps, completedStep, context, sagaId);
                saga.setStatus("COMPENSATED");
                sagaRepo.save(saga);
                return SagaResult.failed(sagaId, step.getName(), e.getMessage());
            }
        }

        saga.setStatus("COMPLETED");
        sagaRepo.save(saga);
        return SagaResult.success(sagaId, context);
    }

    private void compensate(List<SagaStep> steps, int failedAtIndex,
                            Map<String, Object> context, String sagaId) {
        // 从失败步骤向前逐个补偿
        for (int i = failedAtIndex; i >= 0; i--) {
            SagaStep step = steps.get(i);
            if (step.getCompensateAction() == null) continue;
            try {
                log.info("Saga[{}] compensating step: {}", sagaId, step.getName());
                callService(step.getCompensateEndpoint(), step.getCompensateMethod(), context);
            } catch (Exception e) {
                log.error("Saga[{}] compensation step {} failed: {}", sagaId, step.getName(), e.getMessage());
                // 补偿失败不能阻断,记录人工处理
                alertCompensationFailure(sagaId, step.getName(), e.getMessage());
            }
        }
    }

    private Object callService(String endpoint, String method, Map<String, Object> context) {
        // HTTP调用实现,含重试逻辑
        int maxRetries = 3;
        for (int attempt = 1; attempt <= maxRetries; attempt++) {
            try {
                ResponseEntity<Object> response = restTemplate.exchange(
                    endpoint, HttpMethod.resolve(method), new HttpEntity<>(context), Object.class
                );
                return response.getBody();
            } catch (Exception e) {
                if (attempt == maxRetries) throw e;
                try { Thread.sleep(1000 * attempt); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); }
            }
        }
        throw new RuntimeException("Service call failed after retries");
    }
}

补偿操作的设计原则

补偿操作不是”撤销”,而是”业务层面的逆向操作”。设计时必须遵循以下原则:

幂等性:补偿操作必须可安全重试。例如退款接口在已退款时应返回成功而非报错。

// 幂等退款实现
@Service
public class PaymentService {

    private final PaymentRepository paymentRepo;

    public RefundResult refund(String paymentId) {
        Payment payment = paymentRepo.findById(paymentId)
            .orElseThrow(() -> new PaymentNotFoundException(paymentId));

        // 幂等检查:已退款直接返回成功
        if (payment.getStatus() == PaymentStatus.REFUNDED) {
            return RefundResult.alreadyRefunded(paymentId);
        }

        if (payment.getStatus() != PaymentStatus.COMPLETED) {
            throw new InvalidPaymentStateException(paymentId, payment.getStatus());
        }

        payment.setStatus(PaymentStatus.REFUNDED);
        payment.setRefundTime(Instant.now());
        paymentRepo.save(payment);

        return RefundResult.success(paymentId);
    }
}

语义补偿:有些操作无法物理撤销。例如”发送短信”的补偿不是”撤回短信”,而是”发送一条致歉短信”。补偿操作必须从业务语义出发设计,而非追求数据层的逆向操作。

Saga状态持久化与恢复

Saga执行过程中,协调器可能崩溃。重启后需要恢复到崩溃前的状态。方案:每个步骤执行前后持久化Saga实例状态到数据库。

// SagaInstance实体
@Entity
@Table(name = "saga_instances")
public class SagaInstance {
    @Id
    private String sagaId;
    private String status;  // RUNNING, COMPLETED, COMPENSATED, FAILED
    private int currentStep;
    private long createdAt;
    private long updatedAt;

    // JSON存储上下文数据(步骤间传递的中间结果)
    @Column(columnDefinition = "TEXT")
    private String contextData;
}

恢复逻辑:定时扫描status=RUNNING且updatedAt超过5分钟的Saga实例,判断当前步骤是否已完成,未完成则重试或补偿。这一逻辑可以由Spring的@Scheduled任务驱动。

监控与告警指标

Saga运行时必须监控以下指标:Saga成功率(COMPLETED / TOTAL,目标>99%)、平均执行耗时(P50小于2s, P99小于10s)、补偿触发率(COMPENSATED / TOTAL,阈值小于1%)、补偿失败数(需人工介入,告警触发条件大于0)。

建议在SagaOrchestrator中埋点上报Prometheus指标,Grafana面板实时展示。补偿失败必须P1告警,因为这意味着业务数据处于不一致状态。

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

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

相关推荐