分布式事务Saga模式实战:订单系统的一致性保障方案

为什么订单系统需要Saga而非2PC

电商订单流程涉及多个微服务:订单服务创建订单、库存服务扣减库存、支付服务处理付款、物流服务生成运单。传统2PC(两阶段提交)要求所有参与者持有锁直到事务完成,在微服务架构下,锁持有时间可能长达数秒甚至数十秒,高并发场景下直接导致吞吐量崩溃。2PC还有一个致命问题:要求所有参与者实现XA协议,而很多NoSQL数据库和消息队列不支持。

Saga模式的核心思路是将长事务拆分为多个本地短事务,每个本地事务提交后通过消息或事件触发下一个步骤。如果某步失败,执行补偿事务回滚已完成的步骤。这是用最终一致性替代强一致性的工程权衡。

两种Saga编排模式

Saga有两种实现模式:编排式(Choreography)和协调式(Orchestration)。订单场景推荐协调式,因为业务流程清晰可控。

// 编排式:各服务监听事件自行决策(适合简单流程)
// 订单服务发布OrderCreated事件
// 库存服务监听后扣减库存,发布InventoryReserved事件
// 支付服务监听后处理支付,发布PaymentCompleted事件
// 缺点:流程逻辑分散在各服务中,调试困难

// 协调式:由Saga协调器统一驱动(适合复杂流程)
// 协调器按顺序调用各服务,失败时按反序调用补偿

协调式Saga实现:订单创建流程

以一个完整的订单创建Saga为例,展示Java实现:

// Saga定义:定义每个步骤及其补偿操作
public class CreateOrderSaga {

    public static SagaDefinition<CreateOrderState> definition() {
        return SagaDefinition.<CreateOrderState>builder()
            .step()
                .withParticipantName("inventory")
                .withCompensatingAction("reserve")
                .withCompensatingAction("release")
            .step()
                .withParticipantName("payment")
                .withCompensatingAction("charge")
                .withCompensatingAction("refund")
            .step()
                .withParticipantName("logistics")
                .withCompensatingAction("createShipment")
                .withCompensatingAction("cancelShipment")
            .build();
    }
}

// Saga状态机
@Data
public class CreateOrderState {
    private String orderId;
    private String userId;
    private List<OrderItem> items;
    private BigDecimal totalAmount;
    private String reservationId;
    private String paymentId;
    private String shipmentId;
    private SagaStatus status;
    private String currentStep;
    private String failureReason;
}

// Saga协调器核心逻辑
@Service
public class SagaOrchestrator {

    private final SagaStepExecutor stepExecutor;
    private final SagaStateRepository stateRepository;
    private final ApplicationEventPublisher eventPublisher;

    @Transactional
    public CreateOrderState execute(CreateOrderState state) {
        state.setStatus(SagaStatus.EXECUTING);
        stateRepository.save(state);

        try {
            // Step 1: 预留库存
            InventoryResponse invResp = stepExecutor.reserveInventory(
                state.getOrderId(), state.getItems()
            );
            state.setReservationId(invResp.getReservationId());
            state.setCurrentStep("inventory");
            stateRepository.save(state);

            // Step 2: 处理支付
            PaymentResponse payResp = stepExecutor.processPayment(
                state.getOrderId(), state.getUserId(), state.getTotalAmount()
            );
            if (!payResp.isSuccess()) {
                throw new SagaStepException("payment", payResp.getFailureReason());
            }
            state.setPaymentId(payResp.getPaymentId());
            state.setCurrentStep("payment");
            stateRepository.save(state);

            // Step 3: 创建运单
            ShipmentResponse shipResp = stepExecutor.createShipment(
                state.getOrderId(), state.getReservationId()
            );
            state.setShipmentId(shipResp.getShipmentId());
            state.setStatus(SagaStatus.COMPLETED);
            stateRepository.save(state);

            eventPublisher.publishEvent(new SagaCompletedEvent(state));
            return state;

        } catch (SagaStepException e) {
            compensate(state, e.getStepName());
            return state;
        }
    }

    private void compensate(CreateOrderState state, String failedStep) {
        state.setStatus(SagaStatus.COMPENSATING);
        stateRepository.save(state);

        try {
            // 按反序补偿已完成的步骤
            if (failedStep.equals("payment") || 
                List.of("logistics").contains(failedStep)) {
                // 已支付,需要退款
                if (state.getPaymentId() != null) {
                    stepExecutor.refundPayment(state.getPaymentId());
                }
            }
            // 已预留库存,需要释放
            if (state.getReservationId() != null) {
                stepExecutor.releaseInventory(state.getReservationId());
            }
            state.setStatus(SagaStatus.COMPENSATED);
        } catch (Exception e) {
            state.setStatus(SagaStatus.COMPENSATION_FAILED);
            eventPublisher.publishEvent(
                new SagaCompensationFailedEvent(state, e.getMessage())
            );
        }
        stateRepository.save(state);
    }
}

补偿事务的幂等性保证

补偿操作必须幂等——网络超时重试时不能重复退款或重复释放库存。实现幂等的通用方案是使用幂等键:

@Service
public class PaymentService {

    private final IdempotencyKeyRepository idempotencyRepo;

    public RefundResult refund(String paymentId, String idempotencyKey) {
        // 检查幂等键是否已处理
        if (idempotencyRepo.existsByKey(idempotencyKey)) {
            return idempotencyRepo.getResultByKey(idempotencyKey);
        }

        try {
            RefundResult result = callPaymentGateway(paymentId);
            idempotencyRepo.save(idempotencyKey, result);
            return result;
        } catch (Exception e) {
            // 不记录幂等键,允许重试
            throw e;
        }
    }
}

Saga状态持久化与异常恢复

协调器在执行过程中可能崩溃。重启后需要能恢复到中断点继续执行。状态持久化是关键:

@Scheduled(fixedDelay = 30000)
public void recoverStuckSagas() {
    // 查找超时未完成的Saga(超过5分钟仍在执行)
    List<CreateOrderState> stuck = stateRepository
        .findByStatusAndUpdatedAtBefore(
            SagaStatus.EXECUTING,
            LocalDateTime.now().minusMinutes(5)
        );

    stuck.forEach(state -> {
        log.warn("发现卡住的Saga: {}, 当前步骤: {}",
            state.getOrderId(), state.getCurrentStep());
        // 根据当前步骤决定恢复策略
        if (retryCount(state) < 3) {
            retryCurrentStep(state);
        } else {
            compensate(state, state.getCurrentStep());
        }
    });
}

补偿失败是最严重的场景。此时需要人工介入,通过运维面板查看失败原因,手动执行补偿或数据修复。所以Saga状态表中必须保留足够的上下文信息(每步的请求/响应),才能支撑人工修复。

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

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

相关推荐