分布式事务实战:Saga模式在微服务架构中的落地与补偿设计
微服务架构下,一个业务操作跨多个服务的数据修改无法依赖本地事务保证一致性。Saga模式通过将长事务拆分为一系列本地事务,并在失败时执行补偿操作来达到最终一致性。本文以订单支付场景为例,给出Saga模式的具体实现方案和补偿策略设计。
Saga协调模式:编排式 vs 协同式
Saga有两种实现方式:
- 编排式(Orchestration):中央协调器控制事务执行顺序,各服务只执行协调器下发的命令。逻辑集中,流程清晰,适合复杂流程。
- 协同式(Choreography):各服务监听事件并决定自身行为,无中央协调器。服务松耦合,但流程难以追踪,适合简单流程。
生产环境推荐编排式,以下实现基于编排模式。
订单支付Saga流程设计
用户下单支付的完整流程涉及4个服务:订单服务、库存服务、支付服务、通知服务。每个步骤都有对应的补偿操作:
正向流程:
1. 创建订单(Order Service) → 补偿:取消订单
2. 扣减库存(Inventory Service) → 补偿:回滚库存
3. 执行支付(Payment Service) → 补偿:退款
4. 发送通知(Notification Service) → 补偿:无需补偿(幂等)
Saga协调器实现
// Saga协调器核心逻辑
public class OrderPaymentSaga {
private final OrderService orderService;
private final InventoryService inventoryService;
private final PaymentService paymentService;
private final SagaLogRepository sagaLog;
public SagaExecutionResult execute(OrderRequest request) {
String sagaId = UUID.randomUUID().toString();
// 步骤1:创建订单
Order order = orderService.create(request);
sagaLog.save(new SagaStep(sagaId, "CREATE_ORDER", "SUCCESS", order.getId()));
try {
// 步骤2:扣减库存
inventoryService.deduct(order.getProductId(), order.getQuantity());
sagaLog.save(new SagaStep(sagaId, "DEDUCT_INVENTORY", "SUCCESS",
order.getProductId()));
try {
// 步骤3:执行支付
PaymentResult payment = paymentService.charge(order.getId(), order.getAmount());
sagaLog.save(new SagaStep(sagaId, "PAYMENT", "SUCCESS",
payment.getTransactionId()));
// 步骤4:发送通知(非关键步骤,失败不影响整体)
try {
notificationService.sendOrderConfirmation(order);
} catch (Exception e) {
log.warn("Notification failed, saga continues: {}", e.getMessage());
}
return SagaExecutionResult.success(sagaId);
} catch (PaymentException e) {
// 步骤3失败:补偿步骤2和步骤1
compensate(sagaId, "DEDUCT_INVENTORY", () ->
inventoryService.restore(order.getProductId(), order.getQuantity()));
compensate(sagaId, "CREATE_ORDER", () ->
orderService.cancel(order.getId()));
return SagaExecutionResult.failure(sagaId, "PAYMENT_FAILED", e);
}
} catch (InventoryException e) {
// 步骤2失败:补偿步骤1
compensate(sagaId, "CREATE_ORDER", () ->
orderService.cancel(order.getId()));
return SagaExecutionResult.failure(sagaId, "INVENTORY_FAILED", e);
}
}
private void compensate(String sagaId, String stepName, Runnable compensation) {
try {
compensation.run();
sagaLog.save(new SagaStep(sagaId, "COMPENSATE_" + stepName, "SUCCESS", null));
} catch (Exception e) {
sagaLog.save(new SagaStep(sagaId, "COMPENSATE_" + stepName, "FAILED",
e.getMessage()));
// 补偿失败需要人工介入,发送告警
alertService.sendCompensationFailureAlert(sagaId, stepName, e);
}
}
}
Saga日志与恢复机制
Saga协调器崩溃后需要恢复执行中的Saga事务。通过持久化日志记录每一步执行状态,重启后从断点继续:
-- Saga日志表结构
CREATE TABLE saga_log (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
saga_id VARCHAR(64) NOT NULL,
step_name VARCHAR(64) NOT NULL,
status VARCHAR(16) NOT NULL, -- SUCCESS / FAILED / COMPENSATING
payload TEXT,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
INDEX idx_saga_id (saga_id)
);
// 恢复处理器
@Component
public class SagaRecoveryHandler {
@Scheduled(fixedDelay = 30000) // 每30秒扫描
public void recoverPendingSagas() {
List<String> pendingSagas = sagaLog.findIncompleteSagas();
for (String sagaId : pendingSagas) {
SagaState state = sagaLog.getSagaState(sagaId);
if (state.getLastStepStatus().equals("FAILED")) {
// 执行补偿
executeCompensation(sagaId, state);
} else if (state.getLastStepStatus().equals("SUCCESS")) {
// 从下一步继续
continueSaga(sagaId, state);
}
}
}
}
补偿操作的幂等性设计
补偿操作可能被重复执行(网络重试、恢复机制触发),必须保证幂等。通用方案是使用唯一请求ID + 状态机控制:
// 库存服务补偿 - 幂等实现
@Service
public class InventoryService {
@Transactional
public void restore(String productId, int quantity) {
// 1. 检查是否已执行过补偿
CompensationRecord existing = compensationRepo
.findByOperationAndTarget("RESTORE_INVENTORY", productId);
if (existing != null && existing.getStatus().equals("COMPLETED")) {
return; // 已执行,跳过
}
// 2. 执行补偿
inventoryRepo.increaseStock(productId, quantity);
// 3. 记录补偿结果
compensationRepo.save(new CompensationRecord(
"RESTORE_INVENTORY", productId, "COMPLETED"
));
}
@Transactional
public void deduct(String productId, int quantity) {
// 扣减也需幂等:检查当前库存是否足够
int current = inventoryRepo.getStock(productId);
if (current < quantity) {
throw new InsufficientStockException(productId, current, quantity);
}
inventoryRepo.decreaseStock(productId, quantity);
}
}
与TCC模式的对比与选型
TCC(Try-Confirm-Cancel)是另一种分布式事务方案,与Saga的区别在于资源预留阶段:
// TCC实现示例
public interface OrderTccService {
@Transactional
OrderTryResult tryCreate(OrderRequest request) {
// Try阶段:创建待确认订单,预留资源
Order order = new Order();
order.setStatus("PENDING");
orderRepo.save(order);
return new OrderTryResult(order.getId());
}
@Transactional
void confirmCreate(String orderId) {
// Confirm阶段:确认订单,提交事务
Order order = orderRepo.findById(orderId);
order.setStatus("CONFIRMED");
orderRepo.save(order);
}
@Transactional
void cancelCreate(String orderId) {
// Cancel阶段:取消订单,释放预留资源
Order order = orderRepo.findById(orderId);
order.setStatus("CANCELLED");
orderRepo.save(order);
}
}
选型对比:
- 资源占用:TCC在Try阶段预留资源,Confirm前其他事务无法使用;Saga不预留资源,但中间状态对其他事务可见
- 一致性强度:TCC为强一致性(Try成功后Confirm几乎不会失败);Saga为最终一致性
- 开发成本:TCC每个操作需实现3个方法(Try/Confirm/Cancel);Saga只需正向操作和补偿操作
- 适用场景:资金类强一致性要求用TCC;订单、库存等允许最终一致性的用Saga
生产环境注意事项
Saga模式落地时需关注以下工程问题:
- 补偿操作必须比正向操作更稳定——支付服务不可用时退款也会失败,需设计异步重试队列
- Saga事务的隔离性问题——步骤2已扣减库存但步骤3支付尚未完成时,其他请求看到的库存已减少,可能误判为库存不足。解决方案是引入”预扣”状态,查询时区分已扣减和预扣量
- 超时控制——每个本地事务设置超时阈值(如5秒),超时后协调器主动触发补偿
- 监控告警——Saga执行成功率、平均耗时、补偿触发率作为核心指标,补偿触发率高于5%需告警排查
- 死信队列——补偿操作失败超过3次后转入死信队列,人工介入处理
以上方案在日均10万订单的生产环境中验证运行,Saga事务成功率99.7%,补偿触发率0.3%,平均事务耗时1.2秒。补偿失败转入人工处理的比例低于0.01%。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/fen-bu-shi-shi-wu-shi-zhan-saga-mo-shi-zai-wei-fu-wu-jia-2/