分布式事务 Saga 模式实战:构建可靠的数据一致性架构
在微服务架构中,一个业务操作往往需要跨多个服务的数据库进行操作。传统的 ACID 事务在分布式场景下力不从心——两阶段提交(2PC)协议虽然保证了强一致性,但带来了严重的性能瓶颈和可用性问题。Saga 模式应运而生,它通过将长事务拆分为一系列本地事务配合补偿操作,在保证最终一致性的同时,大幅提升了系统的可用性和吞吐量。
本文将深入剖析 Saga 模式的两种核心实现方式(编排与协同),提供完整的代码实战,并分享生产环境中的关键经验与踩坑指南。
1. 为什么微服务需要 Saga?
1.1 分布式事务的现实困境
想象一个电商下单流程:用户点击”提交订单”后,系统需要依次执行:
- 订单服务:创建订单,状态设为”待支付”
- 库存服务:扣减商品库存
- 支付服务:扣减用户余额
- 积分服务:增加用户积分
- 通知服务:发送订单确认通知
这五个操作分布在五个不同的微服务中,每个服务管理自己的数据库。如果使用本地事务,任何一个步骤失败都会导致数据不一致——比如库存扣减了但订单没创建成功,或者支付扣款了但库存没扣减。
1.2 两阶段提交的局限
传统的 2PC 协议通过 Prepare 和 Commit 两个阶段保证原子性,但在微服务架构中面临严峻挑战:
- 同步阻塞:所有参与者在 Prepare 阶段锁定资源,等待协调者决策,长时间持有数据库锁
- 单点故障:协调者宕机导致所有参与者阻塞,整个系统不可用
- 网络分区敏感:网络抖动可能导致事务悬挂,需要人工介入
- 吞吐量低:全局锁导致并发性能急剧下降,通常 TPS 不超过数百
1.3 Saga 的核心理念
Saga 模式的核心思想是:将一个长事务拆分为一系列本地事务(Local Transaction),每个本地事务都有对应的补偿操作(Compensating Transaction)。如果某个步骤失败,系统按逆序执行已成功步骤的补偿操作,回滚到初始状态。
💡 Saga 牺牲了强一致性(Strong Consistency),换来了高可用性(High Availability)和最终一致性(Eventual Consistency)。在大多数业务场景中,这是完全可接受的权衡。
2. Saga 的两种实现模式
2.1 编排模式(Orchestration)
编排模式引入一个中央协调器(Saga Orchestrator),负责按顺序调用各个服务的本地事务,并在失败时触发补偿流程。协调器持有整个 Saga 的状态机,决定下一步执行什么操作。
│ Saga Orchestrator │
│ │
│ ┌─────┐ ┌─────┐ ┌─────┐ ┌─────┐ ┌─────┐ │
│ │Start│───▶│ T1 │───▶│ T2 │───▶│ T3 │───▶│ T4 │ │
│ └─────┘ └──┬──┘ └──┬──┘ └──┬──┘ └──┬──┘ │
│ │ │ │ │ │
│ ┌────▼────┐┌────▼────┐┌────▼────┐┌────▼────┐ │
│ │OrderSvc ││StockSvc ││PaySvc ││PointSvc │ │
│ └─────────┘└─────────┘└─────────┘└─────────┘ │
│ │
│ 失败时: T3 失败 → C2 (补偿T2) → C1 (补偿T1) │
└──────────────────────────────────────────────────────────┘
优点:流程集中管理,易于理解和调试;方便添加监控和重试逻辑;避免循环依赖。
缺点:协调器成为单点;业务逻辑集中在协调器中,可能导致”上帝类”。
2.2 协同模式(Choreography)
协同模式没有中央协调器,每个服务在完成本地事务后发布领域事件(Domain Event),其他服务监听事件并触发自己的操作。整个 Saga 通过事件驱动的方式”自组织”完成。
│ Event Bus (Kafka / RabbitMQ) │
│ │
│ OrderSvc ──OrderCreated──▶ StockSvc │
│ StockSvc ──StockDeducted──▶ PaySvc │
│ PaySvc ──PaymentDone──▶ PointSvc │
│ │
│ 失败时: │
│ PaySvc ──PaymentFailed──▶ StockSvc (补偿: 恢复库存) │
│ StockSvc ──StockRestored──▶ OrderSvc (补偿: 取消订单) │
└──────────────────────────────────────────────────────────┘
优点:去中心化,无单点故障;服务间松耦合;天然支持事件驱动架构。
缺点:流程分散,难以追踪全局状态;调试困难;可能出现循环事件依赖。
2.3 选型决策表
| 考量维度 | 编排模式 | 协同模式 |
|---|---|---|
| 团队规模 | 小团队更易管理 | 大团队职责清晰 |
| 事务复杂度 | ≤7 个步骤的线性流程 | 复杂、动态的流程 |
| 可观测性需求 | 天然支持(集中状态) | 需要额外的事件追踪 |
| 服务耦合度 | 协调器依赖各服务 | 服务间通过事件松耦合 |
| 回滚策略 | 统一由协调器处理 | 各服务自行处理补偿 |
| 推荐场景 | 订单流程、审批流 | 跨部门业务流程、IoT |
3. 编排模式完整实战代码
以下是一个基于 Spring Boot + Kafka 的 Saga 编排模式完整实现,以电商下单流程为例。
3.1 Saga 状态机定义
// Saga 状态枚举
public enum SagaState {
STARTED, // 已启动
ORDER_CREATED, // 订单已创建
STOCK_DEDUCTED, // 库存已扣减
PAYMENT_DONE, // 支付已完成
POINTS_ADDED, // 积分已增加
COMPLETED, // 全部完成
COMPENSATING, // 补偿中
COMPENSATED // 已补偿
}
// Saga 步骤定义
public enum SagaStep {
CREATE_ORDER,
DEDUCT_STOCK,
PROCESS_PAYMENT,
ADD_POINTS,
SEND_NOTIFICATION
}
// Saga 上下文 — 贯穿整个事务生命周期
@Data
public class SagaContext {
private String sagaId;
private SagaState state;
private SagaStep currentStep;
private String orderId;
private String userId;
private String productId;
private Integer quantity;
private BigDecimal amount;
private List<SagaStep> completedSteps = new ArrayList<>();
private String failureReason;
private LocalDateTime createdAt;
private LocalDateTime updatedAt;
}
3.2 Saga 协调器核心实现
@Service
@Slf4j
public class CreateOrderSagaOrchestrator {
private final SagaContextRepository sagaRepository;
private final OrderServiceClient orderService;
private final StockServiceClient stockService;
private final PaymentServiceClient paymentService;
private final PointServiceClient pointService;
private final KafkaTemplate<String, SagaEvent> kafkaTemplate;
/**
* 启动 Saga 事务
*/
@Transactional
public String start(CreateOrderRequest request) {
// 1. 生成 Saga ID 并初始化上下文
String sagaId = UUID.randomUUID().toString();
SagaContext context = new SagaContext();
context.setSagaId(sagaId);
context.setState(SagaState.STARTED);
context.setUserId(request.getUserId());
context.setProductId(request.getProductId());
context.setQuantity(request.getQuantity());
context.setAmount(request.getAmount());
context.setCreatedAt(LocalDateTime.now());
sagaRepository.save(context);
log.info("Saga [{}] 启动: 用户={}, 商品={}, 数量={}",
sagaId, request.getUserId(),
request.getProductId(), request.getQuantity());
// 2. 执行第一步:创建订单
executeStep(context, SagaStep.CREATE_ORDER);
return sagaId;
}
/**
* 执行 Saga 步骤(带重试和超时保护)
*/
private void executeStep(SagaContext context, SagaStep step) {
context.setCurrentStep(step);
context.setUpdatedAt(LocalDateTime.now());
try {
switch (step) {
case CREATE_ORDER -> {
CreateOrderResponse resp = orderService.createOrder(
context.getUserId(),
context.getProductId(),
context.getQuantity()
);
context.setOrderId(resp.getOrderId());
context.setState(SagaState.ORDER_CREATED);
}
case DEDUCT_STOCK -> {
stockService.deductStock(
context.getProductId(),
context.getQuantity()
);
context.setState(SagaState.STOCK_DEDUCTED);
}
case PROCESS_PAYMENT -> {
paymentService.deduct(
context.getUserId(),
context.getAmount()
);
context.setState(SagaState.PAYMENT_DONE);
}
case ADD_POINTS -> {
pointService.addPoints(
context.getUserId(),
context.getAmount()
.multiply(BigDecimal.valueOf(0.01))
);
context.setState(SagaState.POINTS_ADDED);
}
case SEND_NOTIFICATION -> {
// 通知步骤通常可以异步,失败不影响主流程
kafkaTemplate.send("order-notifications",
new OrderNotificationEvent(
context.getOrderId(),
context.getUserId(),
"ORDER_CONFIRMED"
)
);
context.setState(SagaState.COMPLETED);
}
}
// 记录已完成的步骤
context.getCompletedSteps().add(step);
sagaRepository.save(context);
log.info("Saga [{}] 步骤 {} 成功", context.getSagaId(), step);
// 触发下一步
SagaStep nextStep = getNextStep(step);
if (nextStep != null) {
executeStep(context, nextStep);
} else {
// Saga 完成,发布完成事件
publishSagaCompleted(context);
}
} catch (Exception e) {
log.error("Saga [{}] 步骤 {} 失败: {}",
context.getSagaId(), step, e.getMessage());
context.setFailureReason(e.getMessage());
compensate(context);
}
}
/**
* 补偿流程 — 按逆序执行已成功步骤的补偿操作
*/
private void compensate(SagaContext context) {
context.setState(SagaState.COMPENSATING);
sagaRepository.save(context);
log.warn("Saga [{}] 开始补偿, 已完成步骤: {}",
context.getSagaId(), context.getCompletedSteps());
// 逆序遍历已完成的步骤,逐个补偿
List<SagaStep> reversedSteps = new ArrayList<>(context.getCompletedSteps());
Collections.reverse(reversedSteps);
for (SagaStep step : reversedSteps) {
try {
switch (step) {
case CREATE_ORDER ->
orderService.cancelOrder(context.getOrderId());
case DEDUCT_STOCK ->
stockService.restoreStock(
context.getProductId(),
context.getQuantity()
);
case PROCESS_PAYMENT ->
paymentService.refund(
context.getUserId(),
context.getAmount()
);
case ADD_POINTS ->
pointService.deductPoints(
context.getUserId(),
context.getAmount()
.multiply(BigDecimal.valueOf(0.01))
);
}
log.info("Saga [{}] 补偿步骤 {} 成功",
context.getSagaId(), step);
} catch (Exception e) {
// 补偿失败是严重问题,需要告警和人工介入
log.error("Saga [{}] 补偿步骤 {} 失败! 需要人工介入: {}",
context.getSagaId(), step, e.getMessage());
publishCompensationFailed(context, step);
return;
}
}
context.setState(SagaState.COMPENSATED);
sagaRepository.save(context);
log.info("Saga [{}] 补偿完成", context.getSagaId());
}
/**
* 获取下一个步骤
*/
private SagaStep getNextStep(SagaStep current) {
return switch (current) {
case CREATE_ORDER -> SagaStep.DEDUCT_STOCK;
case DEDUCT_STOCK -> SagaStep.PROCESS_PAYMENT;
case PROCESS_PAYMENT -> SagaStep.ADD_POINTS;
case ADD_POINTS -> SagaStep.SEND_NOTIFICATION;
case SEND_NOTIFICATION -> null; // 结束
};
}
}
3.3 幂等性保障 — 补偿操作的核心要求
Saga 模式中,补偿操作可能被多次重试(网络超时、服务重启等),因此每个补偿操作必须是幂等的。以下是库存补偿的幂等实现:
@Service
public class StockCompensationService {
private final StockRepository stockRepository;
private final CompensationLogRepository logRepository;
/**
* 幂等的库存恢复操作
* 通过 compensationId 去重,确保同一补偿只执行一次
*/
@Transactional
public void restoreStock(String compensationId,
String productId,
Integer quantity) {
// 1. 检查是否已执行过
if (logRepository.existsByCompensationId(compensationId)) {
log.info("补偿 [{}] 已执行过,跳过", compensationId);
return;
}
// 2. 使用乐观锁防止并发冲突
int updated = stockRepository.restoreStockWithVersion(
productId, quantity
);
if (updated == 0) {
throw new StockException(
"库存恢复失败: 商品 " + productId + " 版本冲突"
);
}
// 3. 记录补偿日志
CompensationLog logEntry = new CompensationLog();
logEntry.setCompensationId(compensationId);
logEntry.setOperation("RESTORE_STOCK");
logEntry.setProductId(productId);
logEntry.setQuantity(quantity);
logEntry.setExecutedAt(LocalDateTime.now());
logRepository.save(logEntry);
log.info("补偿 [{}] 库存恢复成功: 商品={}, 数量={}",
compensationId, productId, quantity);
}
}
- 唯一标识:每个补偿操作携带全局唯一的 compensationId
- 去重检查:执行前先查询补偿日志表
- 原子记录:补偿执行和日志记录在同一本地事务中
3.4 持久化与恢复 — 协调器宕机怎么办?
编排模式的一个关键风险是协调器本身可能宕机。我们需要确保 Saga 状态被持久化,并在协调器恢复后能继续执行。
@Service
public class SagaRecoveryService {
private final SagaContextRepository sagaRepository;
private final CreateOrderSagaOrchestrator orchestrator;
/**
* 定时任务:每 30 秒扫描未完成的 Saga
* 适用于协调器重启后的恢复
*/
@Scheduled(fixedDelay = 30_000)
public void recoverIncompleteSagas() {
List<SagaContext> incompleteSagas = sagaRepository
.findByStateInAndUpdatedAtBefore(
List.of(
SagaState.STARTED,
SagaState.ORDER_CREATED,
SagaState.STOCK_DEDUCTED,
SagaState.PAYMENT_DONE,
SagaState.POINTS_ADDED,
SagaState.COMPENSATING
),
LocalDateTime.now().minusMinutes(5)
);
for (SagaContext saga : incompleteSagas) {
log.info("恢复 Saga [{}], 当前状态: {}, 当前步骤: {}",
saga.getSagaId(), saga.getState(), saga.getCurrentStep());
try {
if (saga.getState() == SagaState.COMPENSATING) {
// 补偿中的 Saga 继续补偿
orchestrator.resumeCompensation(saga);
} else {
// 执行中的 Saga 从当前步骤继续
orchestrator.resumeFromStep(
saga, saga.getCurrentStep()
);
}
} catch (Exception e) {
log.error("恢复 Saga [{}] 失败: {}",
saga.getSagaId(), e.getMessage());
}
}
}
}
4. 协同模式实战代码
以下展示协同模式的核心实现,使用 Kafka 作为事件总线。
4.1 事件定义
// 基础事件接口
public interface SagaEvent {
String getSagaId();
String getEventType();
}
// 订单创建事件
@Data
@AllArgsConstructor
public class OrderCreatedEvent implements SagaEvent {
private String sagaId;
private String orderId;
private String userId;
private String productId;
private Integer quantity;
private BigDecimal amount;
}
// 库存扣减成功事件
@Data
@AllArgsConstructor
public class StockDeductedEvent implements SagaEvent {
private String sagaId;
private String productId;
private Integer quantity;
}
// 支付失败事件 — 触发补偿链
@Data
@AllArgsConstructor
public class PaymentFailedEvent implements SagaEvent {
private String sagaId;
private String orderId;
private String userId;
private String productId;
private Integer quantity;
private String reason;
}
4.2 事件处理器 — 每个服务独立处理
@Service
@Slf4j
public class StockSagaHandler {
private final StockService stockService;
private final KafkaTemplate<String, SagaEvent> kafkaTemplate;
/**
* 监听 OrderCreated 事件 → 扣减库存
*/
@KafkaListener(
topics = "saga-events",
groupId = "stock-service",
containerFactory = "sagaEventListenerFactory"
)
public void handleOrderCreated(OrderCreatedEvent event) {
log.info("收到 OrderCreated 事件, Saga [{}]", event.getSagaId());
try {
stockService.deductStock(
event.getProductId(), event.getQuantity()
);
// 扣减成功 → 发布 StockDeducted 事件
kafkaTemplate.send("saga-events",
new StockDeductedEvent(
event.getSagaId(),
event.getProductId(),
event.getQuantity()
)
);
} catch (Exception e) {
log.error("库存扣减失败, Saga [{}]: {}",
event.getSagaId(), e.getMessage());
// 扣减失败 → 发布 StockDeductFailed 事件(触发补偿)
kafkaTemplate.send("saga-events",
new StockDeductFailedEvent(
event.getSagaId(),
event.getOrderId(),
event.getProductId(),
event.getQuantity(),
e.getMessage()
)
);
}
}
/**
* 监听 PaymentFailed 事件 → 补偿:恢复库存
*/
@KafkaListener(
topics = "saga-events",
groupId = "stock-service"
)
public void handlePaymentFailed(PaymentFailedEvent event) {
log.info("收到 PaymentFailed 事件, Saga [{}] 开始补偿",
event.getSagaId());
try {
stockService.restoreStock(
event.getProductId(), event.getQuantity()
);
// 补偿完成 → 发布 StockRestored 事件
kafkaTemplate.send("saga-events",
new StockRestoredEvent(
event.getSagaId(),
event.getProductId(),
event.getQuantity()
)
);
} catch (Exception e) {
log.error("库存恢复补偿失败, Saga [{}]: {}",
event.getSagaId(), e.getMessage());
// 发送告警,需要人工介入
}
}
}
5. 生产级最佳实践
5.1 超时与死锁防护
Saga 中的每个步骤都应该设置合理的超时时间。如果一个步骤长时间未完成,应该触发补偿流程。
@Configuration
public class SagaTimeoutConfig {
/**
* Saga 步骤超时配置
* 不同步骤的超时时间根据业务特性设置
*/
@Bean
public Map<SagaStep, Duration> sagaStepTimeouts() {
return Map.of(
SagaStep.CREATE_ORDER, Duration.ofSeconds(10),
SagaStep.DEDUCT_STOCK, Duration.ofSeconds(5),
SagaStep.PROCESS_PAYMENT, Duration.ofSeconds(30), // 支付涉及第三方
SagaStep.ADD_POINTS, Duration.ofSeconds(5),
SagaStep.SEND_NOTIFICATION, Duration.ofSeconds(3)
);
}
/**
* Saga 整体超时 — 超过此时间未完成则标记为失败
*/
@Bean
public Duration sagaGlobalTimeout() {
return Duration.ofMinutes(5);
}
}
5.2 可观测性 — 追踪与监控
Saga 的分布式特性使得可观测性至关重要。建议集成 OpenTelemetry 进行全链路追踪:
@Aspect
@Component
public class SagaObservabilityAspect {
private final MeterRegistry meterRegistry;
private final Tracer tracer;
@Around("@annotation(SagaStep)")
public Object observeSagaStep(ProceedingJoinPoint pjp)
throws Throwable {
String stepName = pjp.getSignature().getName();
String sagaId = extractSagaId(pjp.getArgs());
Span span = tracer.spanBuilder("saga-step-" + stepName)
.setAttribute("saga.id", sagaId)
.setAttribute("saga.step", stepName)
.startSpan();
Timer.Sample sample = Timer.start(meterRegistry);
try (Scope scope = span.makeCurrent()) {
Object result = pjp.proceed();
span.setAttribute("saga.success", true);
meterRegistry.counter("saga.step.success",
"step", stepName).increment();
return result;
} catch (Exception e) {
span.setAttribute("saga.success", false);
span.setAttribute("error.message", e.getMessage());
meterRegistry.counter("saga.step.failure",
"step", stepName).increment();
throw e;
} finally {
sample.stop(meterRegistry.timer("saga.step.duration",
"step", stepName));
span.end();
}
}
}
5.3 关键经验总结
- 补偿操作必须幂等:网络超时可能导致补偿被重试,非幂等操作会造成数据错误
- 补偿可能失败:补偿操作不是 100% 可靠的,需要设计告警机制和人工介入流程
- 避免 Saga 嵌套:Saga 步骤中不要嵌套另一个 Saga,否则复杂度会指数级增长
- 数据隔离性:Saga 的中间状态对其他事务不可见,考虑使用语义锁(Semantic Lock)
- 版本兼容:补偿逻辑需要兼容旧版本的数据格式,避免升级时补偿失败
6. 2026 年趋势展望
Saga 模式在云原生时代正在持续演进,以下是值得关注的趋势:
- Saga + Event Sourcing 融合:将 Saga 的状态变更以事件溯源的方式持久化,天然支持审计和回放
- AI 驱动的补偿策略:利用机器学习预测补偿失败概率,动态调整重试策略和超时时间
- Serverless Saga:基于 AWS Step Functions / Azure Durable Functions 的无服务器 Saga 编排,按需付费
- Saga 模式语言(Saga DSL):声明式定义 Saga 流程,自动生成协调器和补偿代码
- 与 TCC/3PC 混合使用:对强一致性要求高的步骤使用 TCC,其他步骤使用 Saga,平衡一致性和性能
7. 总结
Saga 模式是微服务架构中解决分布式事务问题的核心方案。通过本文的实战代码和分析,我们总结出以下关键要点:
- 编排模式适合流程简单、团队规模小的场景,集中管理易于调试
- 协同模式适合复杂跨部门流程,去中心化架构更灵活
- 幂等性是补偿操作的基石,必须通过唯一 ID + 去重日志保证
- 持久化与恢复机制确保协调器宕机后 Saga 能自动恢复
- 可观测性是生产环境的必备能力,全链路追踪 + 指标监控缺一不可
在分布式系统中,没有银弹。Saga 模式通过牺牲强一致性换取了可用性和性能,理解这个权衡,才能在正确的场景做出正确的选择。
📖 延伸阅读推荐:
• 《Microservices Patterns》— Chris Richardson(Saga 模式原创作者)
• 《Designing Data-Intensive Applications》— Martin Kleppmann
• Saga Pattern 原始论文:Sagas (1987) — Hector Garcia-Molina & Kenneth Salem