为什么订单系统需要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/