微服务分布式事务实战:Saga模式从理论到Spring Boot落地

为什么不用XA两阶段提交

分布式事务最直觉的方案是XA两阶段提交。协调者通知所有参与者prepare,全部确认后再commit,任何一个拒绝就rollback。理论上完美,工程上几乎不可用——数据库长时间持有行锁,并发性能直接归零;协调者单点故障导致所有参与者阻塞;跨服务部署时网络分区让超时成为常态。

高并发场景下,XA的锁等待时间会从毫秒级膨胀到秒级。一个电商下单流程涉及库存、订单、支付三个服务,XA事务下库存行的锁持有时间覆盖整个流程,高峰期这种设计会直接拖垮库存服务。

Saga模式的思路完全不同:把长事务拆成多个本地短事务,每个本地事务提交后释放锁,通过补偿事务处理失败回滚。锁持有时间从整个分布式事务缩短到单个本地事务,性能提升数量级。

编排式Saga vs 协调式Saga

两种实现流派。编排式(Choreography)是事件驱动,每个服务完成本地事务后发事件通知下游,下游监听事件执行自己的事务。协调式(Orchestration)有一个中央协调器,按照预设流程依次调用各服务。

编排式优点是松耦合,缺点是业务流程散落在各服务的事件监听器里,看代码无法直观理解完整流程。出了问题排查像在拼图。

协调式优点是流程一目了然,缺点是协调器可能成为单点。对于大多数业务系统,协调式的可维护性优势远大于单点风险——协调器本身可以用数据库事务保证可靠性。

生产环境推荐协调式。这篇文章也以协调式Saga为主展开。

Saga状态机设计

Saga的核心是状态机。每个Saga实例有一个状态,事件驱动状态转换:

public enum OrderSagaState {
    ORDER_CREATED,
    STOCK_RESERVED,
    PAYMENT_PROCESSING,
    PAYMENT_COMPLETED,
    COMPLETED,
    STOCK_RESERVATION_FAILED,
    PAYMENT_FAILED,
    COMPENSATING_STOCK,
    COMPENSATED
}

@Data
@Entity
public class OrderSaga {
    @Id
    private String sagaId;
    private String orderId;
    private OrderSagaState state;
    private String stockReservationId;
    private String paymentId;
    private Instant createdAt;
    private Instant updatedAt;
}

状态机定义了正向流程和补偿流程的完整路径。正向路径是ORDER_CREATED -> STOCK_RESERVED -> PAYMENT_PROCESSING -> PAYMENT_COMPLETED -> COMPLETED。任何步骤失败,进入补偿路径反向执行已完成的补偿操作。

Spring Boot实现Saga协调器

核心是SagaManager,负责推进Saga实例的状态转换:

@Service
@Slf4j
public class OrderSagaManager {

    private final OrderSagaRepository sagaRepository;
    private final StockServiceClient stockService;
    private final PaymentServiceClient paymentService;
    private final ApplicationEventPublisher eventPublisher;

    public void handle(OrderCreatedEvent event) {
        OrderSaga saga = new OrderSaga();
        saga.setSagaId(UUID.randomUUID().toString());
        saga.setOrderId(event.getOrderId());
        saga.setState(OrderSagaState.ORDER_CREATED);
        saga.setCreatedAt(Instant.now());
        sagaRepository.save(saga);

        tryReserveStock(saga, event);
    }

    private void tryReserveStock(OrderSaga saga, OrderCreatedEvent event) {
        try {
            String reservationId = stockService.reserve(
                event.getProductId(), event.getQuantity()
            );
            saga.setStockReservationId(reservationId);
            saga.setState(OrderSagaState.STOCK_RESERVED);
            saga.setUpdatedAt(Instant.now());
            sagaRepository.save(saga);

            tryProcessPayment(saga, event);
        } catch (Exception e) {
            log.error("库存预留失败, sagaId={}", saga.getSagaId(), e);
            saga.setState(OrderSagaState.STOCK_RESERVATION_FAILED);
            saga.setUpdatedAt(Instant.now());
            sagaRepository.save(saga);
        }
    }

    private void tryProcessPayment(OrderSaga saga, OrderCreatedEvent event) {
        saga.setState(OrderSagaState.PAYMENT_PROCESSING);
        sagaRepository.save(saga);

        try {
            String paymentId = paymentService.charge(
                event.getUserId(), event.getAmount()
            );
            saga.setPaymentId(paymentId);
            saga.setState(OrderSagaState.PAYMENT_COMPLETED);
            saga.setUpdatedAt(Instant.now());
            sagaRepository.save(saga);

            eventPublisher.publishEvent(
                new OrderCompletedEvent(saga.getOrderId())
            );
        } catch (Exception e) {
            log.error("支付失败, sagaId={}", saga.getSagaId(), e);
            compensateStock(saga);
        }
    }

    private void compensateStock(OrderSaga saga) {
        saga.setState(OrderSagaState.COMPENSATING_STOCK);
        sagaRepository.save(saga);

        try {
            stockService.release(saga.getStockReservationId());
            saga.setState(OrderSagaState.COMPENSATED);
        } catch (Exception e) {
            log.error("库存补偿失败, 需人工介入! sagaId={}", saga.getSagaId(), e);
        }
        saga.setUpdatedAt(Instant.now());
        sagaRepository.save(saga);
    }
}

每个步骤的本地事务由各自服务保证。Saga协调器只负责编排调用顺序和记录状态,不跨服务持锁。这就是Saga比XA快的核心原因。

补偿操作的幂等性保证

补偿操作必须幂等。网络超时情况下,协调器可能重复发起补偿请求。如果补偿操作不是幂等的,比如库存释放两次就会多加库存。

实现幂等的简单方案——唯一请求ID加数据库唯一约束:

@Service
public class StockService {

    private final StockRepository stockRepo;
    private final CompensationLogRepository logRepo;

    @Transactional
    public void releaseReservation(String reservationId, String idempotencyKey) {
        // 幂等检查
        if (logRepo.existsByIdempotencyKey(idempotencyKey)) {
            log.info("补偿操作已执行, 忽略重复请求: {}", idempotencyKey);
            return;
        }

        StockReservation reservation = stockRepo
            .findByReservationId(reservationId)
            .orElseThrow(() -> new ReservationNotFoundException(reservationId));

        // 恢复库存
        stockRepo.addStock(reservation.getProductId(), reservation.getQuantity());

        // 记录补偿日志
        logRepo.save(new CompensationLog(idempotencyKey, reservationId, Instant.now()));
    }
}

IdempotencyKey由协调器生成,通常用sagaId + step名组合。补偿日志表上加唯一约束,并发情况下第二次插入直接失败,事务回滚,达到幂等效果。

Saga异常处理的工程实践

补偿操作本身也可能失败。支付成功但库存释放失败,数据就不一致了。处理这种场景的策略:

1. 重试机制:补偿操作加指数退避重试。大多数临时故障(网络抖动、数据库连接池满)都能通过重试解决。

2. 死信队列:重试超过N次后,将失败记录写入死信队列,人工处理。

3. 定时巡检:后台任务扫描超过一定时间仍处于中间状态的Saga实例,触发补偿或告警。

@Scheduled(fixedRate = 60000)
public void scanStuckSagas() {
    Instant threshold = Instant.now().minus(Duration.ofMinutes(10));
    List<OrderSaga> stuckSagas = sagaRepository
        .findByStateNotInAndUpdatedAtBefore(
            List.of(COMPLETED, COMPENSATED, STOCK_RESERVATION_FAILED),
            threshold
        );

    for (OrderSaga saga : stuckSagas) {
        log.warn("发现停滞Saga: sagaId={}, state={}, 触发补偿",
            saga.getSagaId(), saga.getState());
        compensateFromState(saga);
    }
}

微服务下的分布式事务没有银弹。Saga用最终一致性替代了强一致性,用补偿操作替代了全局回滚。对于绝大多数互联网业务——电商、支付、物流——这种取舍是合理的。理解业务对一致性的真实需求,比选择技术方案更重要。

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

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

相关推荐

微服务分布式事务实战:Saga模式从理论到Spring Boot落地

为什么不用XA两阶段提交

分布式事务最直觉的方案是XA两阶段提交。协调者通知所有参与者prepare,全部确认后再commit,任何一个拒绝就rollback。理论上完美,工程上几乎不可用——数据库长时间持有行锁,并发性能直接归零;协调者单点故障导致所有参与者阻塞;跨服务部署时网络分区让超时成为常态。

高并发场景下,XA的锁等待时间会从毫秒级膨胀到秒级。一个电商下单流程涉及库存、订单、支付三个服务,XA事务下库存行的锁持有时间覆盖整个流程,高峰期这种设计会直接拖垮库存服务。

Saga模式的思路完全不同:把长事务拆成多个本地短事务,每个本地事务提交后释放锁,通过补偿事务处理失败回滚。锁持有时间从整个分布式事务缩短到单个本地事务,性能提升数量级。

编排式Saga vs 协调式Saga

两种实现流派。编排式(Choreography)是事件驱动,每个服务完成本地事务后发事件通知下游,下游监听事件执行自己的事务。协调式(Orchestration)有一个中央协调器,按照预设流程依次调用各服务。

编排式优点是松耦合,缺点是业务流程散落在各服务的事件监听器里,看代码无法直观理解完整流程。出了问题排查像在拼图。

协调式优点是流程一目了然,缺点是协调器可能成为单点。对于大多数业务系统,协调式的可维护性优势远大于单点风险——协调器本身可以用数据库事务保证可靠性。

生产环境推荐协调式。这篇文章也以协调式Saga为主展开。

Saga状态机设计

Saga的核心是状态机。每个Saga实例有一个状态,事件驱动状态转换:

public enum OrderSagaState {
    ORDER_CREATED,
    STOCK_RESERVED,
    PAYMENT_PROCESSING,
    PAYMENT_COMPLETED,
    COMPLETED,
    STOCK_RESERVATION_FAILED,
    PAYMENT_FAILED,
    COMPENSATING_STOCK,
    COMPENSATED
}

@Data
@Entity
public class OrderSaga {
    @Id
    private String sagaId;
    private String orderId;
    private OrderSagaState state;
    private String stockReservationId;
    private String paymentId;
    private Instant createdAt;
    private Instant updatedAt;
}

状态机定义了正向流程和补偿流程的完整路径。正向路径是ORDER_CREATED -> STOCK_RESERVED -> PAYMENT_PROCESSING -> PAYMENT_COMPLETED -> COMPLETED。任何步骤失败,进入补偿路径反向执行已完成的补偿操作。

Spring Boot实现Saga协调器

核心是SagaManager,负责推进Saga实例的状态转换:

@Service
@Slf4j
public class OrderSagaManager {

    private final OrderSagaRepository sagaRepository;
    private final StockServiceClient stockService;
    private final PaymentServiceClient paymentService;
    private final ApplicationEventPublisher eventPublisher;

    public void handle(OrderCreatedEvent event) {
        OrderSaga saga = new OrderSaga();
        saga.setSagaId(UUID.randomUUID().toString());
        saga.setOrderId(event.getOrderId());
        saga.setState(OrderSagaState.ORDER_CREATED);
        saga.setCreatedAt(Instant.now());
        sagaRepository.save(saga);

        tryReserveStock(saga, event);
    }

    private void tryReserveStock(OrderSaga saga, OrderCreatedEvent event) {
        try {
            String reservationId = stockService.reserve(
                event.getProductId(), event.getQuantity()
            );
            saga.setStockReservationId(reservationId);
            saga.setState(OrderSagaState.STOCK_RESERVED);
            saga.setUpdatedAt(Instant.now());
            sagaRepository.save(saga);

            tryProcessPayment(saga, event);
        } catch (Exception e) {
            log.error("库存预留失败, sagaId={}", saga.getSagaId(), e);
            saga.setState(OrderSagaState.STOCK_RESERVATION_FAILED);
            saga.setUpdatedAt(Instant.now());
            sagaRepository.save(saga);
        }
    }

    private void tryProcessPayment(OrderSaga saga, OrderCreatedEvent event) {
        saga.setState(OrderSagaState.PAYMENT_PROCESSING);
        sagaRepository.save(saga);

        try {
            String paymentId = paymentService.charge(
                event.getUserId(), event.getAmount()
            );
            saga.setPaymentId(paymentId);
            saga.setState(OrderSagaState.PAYMENT_COMPLETED);
            saga.setUpdatedAt(Instant.now());
            sagaRepository.save(saga);

            eventPublisher.publishEvent(
                new OrderCompletedEvent(saga.getOrderId())
            );
        } catch (Exception e) {
            log.error("支付失败, sagaId={}", saga.getSagaId(), e);
            compensateStock(saga);
        }
    }

    private void compensateStock(OrderSaga saga) {
        saga.setState(OrderSagaState.COMPENSATING_STOCK);
        sagaRepository.save(saga);

        try {
            stockService.release(saga.getStockReservationId());
            saga.setState(OrderSagaState.COMPENSATED);
        } catch (Exception e) {
            log.error("库存补偿失败, 需人工介入! sagaId={}", saga.getSagaId(), e);
        }
        saga.setUpdatedAt(Instant.now());
        sagaRepository.save(saga);
    }
}

每个步骤的本地事务由各自服务保证。Saga协调器只负责编排调用顺序和记录状态,不跨服务持锁。这就是Saga比XA快的核心原因。

补偿操作的幂等性保证

补偿操作必须幂等。网络超时情况下,协调器可能重复发起补偿请求。如果补偿操作不是幂等的,比如库存释放两次就会多加库存。

实现幂等的简单方案——唯一请求ID加数据库唯一约束:

@Service
public class StockService {

    private final StockRepository stockRepo;
    private final CompensationLogRepository logRepo;

    @Transactional
    public void releaseReservation(String reservationId, String idempotencyKey) {
        // 幂等检查
        if (logRepo.existsByIdempotencyKey(idempotencyKey)) {
            log.info("补偿操作已执行, 忽略重复请求: {}", idempotencyKey);
            return;
        }

        StockReservation reservation = stockRepo
            .findByReservationId(reservationId)
            .orElseThrow(() -> new ReservationNotFoundException(reservationId));

        // 恢复库存
        stockRepo.addStock(reservation.getProductId(), reservation.getQuantity());

        // 记录补偿日志
        logRepo.save(new CompensationLog(idempotencyKey, reservationId, Instant.now()));
    }
}

IdempotencyKey由协调器生成,通常用sagaId + step名组合。补偿日志表上加唯一约束,并发情况下第二次插入直接失败,事务回滚,达到幂等效果。

Saga异常处理的工程实践

补偿操作本身也可能失败。支付成功但库存释放失败,数据就不一致了。处理这种场景的策略:

1. 重试机制:补偿操作加指数退避重试。大多数临时故障(网络抖动、数据库连接池满)都能通过重试解决。

2. 死信队列:重试超过N次后,将失败记录写入死信队列,人工处理。

3. 定时巡检:后台任务扫描超过一定时间仍处于中间状态的Saga实例,触发补偿或告警。

@Scheduled(fixedRate = 60000)
public void scanStuckSagas() {
    Instant threshold = Instant.now().minus(Duration.ofMinutes(10));
    List<OrderSaga> stuckSagas = sagaRepository
        .findByStateNotInAndUpdatedAtBefore(
            List.of(COMPLETED, COMPENSATED, STOCK_RESERVATION_FAILED),
            threshold
        );

    for (OrderSaga saga : stuckSagas) {
        log.warn("发现停滞Saga: sagaId={}, state={}, 触发补偿",
            saga.getSagaId(), saga.getState());
        compensateFromState(saga);
    }
}

微服务下的分布式事务没有银弹。Saga用最终一致性替代了强一致性,用补偿操作替代了全局回滚。对于绝大多数互联网业务——电商、支付、物流——这种取舍是合理的。理解业务对一致性的真实需求,比选择技术方案更重要。

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

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

相关推荐