分布式事务的核心矛盾与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/