Saga 模式:微服务架构中的分布式事务
在传统的整体应用程序中,维护多个实体之间的数据一致性非常简单。关系数据库引擎提供封装在本地 SQL 事务中的 ACID(原子性、一致性、隔离性、持久性)保证。如果下单、付款扣除或库存储备中途失败,调用 ROLLBACK 会立即恢复每个数据库修改。
然而,当迁移到现代微服务架构时,数据管理发生了根本性的变化。为了保证域自治和独立的可扩展性,每个微服务都拥有自己的私有数据库。单个业务操作(例如处理电子商务结帐)现在跨越多个服务边界和数据库引擎(例如,用于订单的 PostgreSQL、用于支付的 DynamoDB、用于库存的 Redis)。
由于分布式微服务不能依赖于单个数据库事务,因此维护跨网络边界的数据一致性成为分布式系统工程中最具挑战性的问题之一。
为了在不牺牲系统可用性或性能的情况下解决这个问题,软件架构师依赖 Saga 模式。
在本次深入研究中,我们将探讨传统分布式事务失败的原因,分解 Saga 模式的核心机制,比较 Choreography 与 Orchestration,分析隔离对策,检查 Go 和 Java 中的生产代码实现,并学习如何安全地处理现实世界的故障回滚。
根本问题:为什么微服务中的两阶段提交 (2PC) 失败
在采用 Saga 模式之前,工程师经常会问:为什么我们不能在微服务中使用传统的两阶段提交(2PC / XA)?
2PC 的机制
两阶段提交使用中央事务管理器分两个阶段协调跨多个数据库节点的分布式事务:
- 准备阶段:事务管理器要求所有参与的数据库节点准备并锁定所需的行。参与者投票
YES或NO。 - 提交阶段:如果所有参与者都投票了
YES,则管理器向所有节点发出COMMIT命令。如果任何节点投票NO或超时,它会发出ROLLBACK。
Client ----> Transaction Manager
|
+-------------+-------------+
| (Prepare) | (Prepare) | (Prepare)
v v v
Order DB Payment DB Inventory DB
为什么 2PC 是微服务的反模式
虽然 2PC 保证了强一致性,但由于几个架构缺陷,它在云原生微服务环境中崩溃了:
- 阻塞和资源争用:在多阶段网络握手过程中,数据库行保持锁定状态。如果网络延迟增加或服务变慢,锁就会保持打开状态,从而快速消耗连接池并导致级联系统故障。
- 可用性瓶颈:在 2PC 中,系统可用性受到所有参与者可用性的乘积的限制($A_{total} = A_1 \times A_2 \times \dots \times A_n$)。如果任何单个服务或数据库节点在准备阶段离线,整个全局事务将无限期阻塞。
- 失去服务自主性:2PC 强制服务跨网络 API 公开数据库级 XA 协议,直接耦合其数据库引擎。
- 代理约束:像 Apache Kafka 或 RabbitMQ 这样的高吞吐量消息代理本身并不参与跨关系数据库的传统 XA 2PC 事务。
根据CAP定理,分布式系统必须在网络分区(P)下的强一致性(C)和高可用性(A)之间进行选择。现代微服务系统优先考虑可用性和分区容错性,用即时一致性换取最终一致性(BASE:基本可用、软状态、最终一致性)。
##什么是传奇模式?
Saga 模式最初由 Hector Garcia-Molina 和 Kenneth Salem 于 1987 年提出,作为处理数据库管理系统中长期事务的机制。在现代微服务中,Saga 代表一系列离散的本地事务。
Saga 不是将整个多服务流包装在单个全局锁内,而是将业务流程作为一系列独立的本地数据库事务($T_1、T_2、\dots、T_n$)来执行。
每个本地事务都会更新单个微服务数据库中的数据并发布域事件或消息。此事件触发下游服务中的下一个本地事务 ($T_{i+1}$)。
[Order Service] [Payment Service] [Inventory Service]
Local Tx T1 -------------> Local Tx T2 -------------> Local Tx T3
(Create Order) (Process Payment) (Reserve Stock)
前向执行与后向补偿
如果所有本地事务都成功,则 Saga 成功完成(转发执行)。
但是,如果本地事务中途失败(例如,$T_3$ 由于商品缺货而失败),则 Saga 不能简单地为之前的步骤($T_1、T_2$)调用数据库 ROLLBACK,因为这些本地事务已经提交到各自的数据库。
要撤消已提交的本地事务,Saga 必须以相反的顺序执行 补偿事务 ($C_{n-1}, \dots, C_1$)。
Forward Flow: T1 (Create Order) ---> T2 (Charge Card) ---> T3 (Reserve Inventory - FAILS)
|
Backward Rollback: C1 (Cancel Order) <--- C2 (Refund Card) <----------+
Saga交易的分类
要设计强大的 Saga 工作流程,事务序列中的每个步骤都必须分为以下三种结构类型之一:
| 交易类型 | 描述 | 幂等性和回滚要求 |
|---|---|---|
| 可补偿交易 | 在无法返回的点之前执行的步骤。如果下游步骤失败,则可以撤消或撤销它们。 | 必须有相应的补偿交易($C_i$)。 |
| 枢轴交易 | 传奇中的决定性一步。如果枢轴成功,传奇就一定会结束。如果失败,Saga 就会回滚。 | 既不可赔偿也不可重审;它标志着回滚和前进完成之间的界限。 |
| 可重审交易 | Pivot 事务之后执行的步骤。他们保证最终会成功,并且不需要补偿。 | 必须严格幂等,因为它们将自动重试直到成功。 |
电子商务结帐示例细分
考虑一个由四个步骤组成的电子商务订单结帐:
- $T_1$:创建挂单(可补偿)$\to$ 通过 $C_1$ 撤消:取消订单。
- $T_2$:授权付款(可补偿)$\to$ 通过 $C_2$ 撤消:退款付款。
- $T_3$: 储备库存 (枢轴交易) $\to$ 如果库存分配成功,则订单最终确定。如果失败,则触发$C_2$和$C_1$。
- $T_4$:派送运输请求 (可重审) $\to$ 在枢轴后执行;重试直至交付。
Saga 架构风格:编排与编排
在分布式系统中实现 Saga 模式有两种主要的架构风格:编排(分散式)和 编排(集中式)。
风格 1:编排(事件驱动的去中心化)
在 基于编排的 Saga 中,没有中央控制器或协调器。相反,微服务通过侦听发布到中央事件总线(例如 Apache Kafka、NATS 或 RabbitMQ)的域事件来进行异步通信。
工作流程机制
- 订单服务 执行 $T_1$ (创建挂单)并向 Kafka 发出
OrderCreated事件。 - 支付服务 侦听
OrderCreated,执行 $T_2$(向信用卡收费),并发出PaymentCompleted事件(或PaymentFailed)。 - 库存服务 监听
PaymentCompleted,执行 $T_3$(保留商品)。如果商品缺货,则会发出InventoryReservationFailed。 - 支付服务 监听
InventoryReservationFailed并执行 $C_2$ (发出退款)。 - 订单服务 监听
PaymentRefunded并执行 $C_1$ (将订单标记为已取消)。
####编排的优点
- 松耦合:服务仅订阅事件主题;他们不知道其他服务的实现。
- 高吞吐量和去中心化:直接事件发布/订阅消除了中央编排器瓶颈。
- 短工作流程的简单性:易于设置简单的 2-3 步工作流程。
编排的缺点
- 循环依赖风险:服务最终可能会监听彼此的事件,从而创建复杂的循环依赖。
- 困难的代码可追溯性:了解端到端业务流程需要跨多个代码库存储库跟踪逻辑。
- 事件风暴和复杂性:随着步骤数量的增加(例如,10 多个服务),管理错误边缘情况会导致事件状态爆炸。
风格 2:编排(集中式工作流程管理)
在 基于编排的 Saga 中,称为 Saga Orchestrator 的专用微服务控制分布式事务的整个生命周期。协调器充当中央协调器,向参与的微服务发出显式命令并监听它们的响应事件。
工作流程机制
- 客户端向 Saga Orchestrator 发送订单请求。
- Orchestrator 向 Order Service 发送
CreateOrder命令。订单服务返回OrderCreated。 - Orchestrator 更新其状态机并向 支付服务 发送
ProcessPayment命令。付款服务返回PaymentSuccessful。 - Orchestrator 向 Inventory Service 发送
ReserveInventory命令。库存服务返回InventoryFailed (Out of Stock)。 - Orchestrator 检测到故障,启动补偿流程:
- 发送
RefundPayment命令至 支付服务。 - 发送
CancelOrder命令至 订单服务。
- 发送
- Orchestrator 将 Saga 执行标记为
FAILED。
编排的优点
- 集中式业务逻辑:工作流状态和业务逻辑本地化在单个协调器服务或状态机中。
- 无循环依赖:微服务响应来自协调器的命令;他们不依赖或了解其他下游服务。
- 清晰的监控和调试:端到端事务状态显式存储在编排器的状态存储中(例如,PostgreSQL 或 Temporal / Camunda 等工作流引擎)。
- 更轻松的错误处理:添加新步骤或更改回滚规则完全在协调器内部进行管理。
编排的缺点
- 编排器复杂性:将太多域逻辑放入编排器中的风险,将其变成“智能编排器,愚蠢的服务”反模式。
- 潜在的单点故障:编排器必须具有高可用性和有状态。
比较矩阵:编排与编排
| 特色 | 编舞 | 编排 |
|---|---|---|
| 控制结构 | 去中心化(事件发布/订阅) | 集中式(Saga 协调器/状态机) |
| 联轴器 | 极低(服务消耗事件) | 中(服务接受协调器的命令) |
| 流程可见性 | 低(跨日志文件分布) | 高(单状态存储可视化工作流程) |
| 最适合 | 简单的工作流程(2 至 4 个服务步骤) | 复杂的企业工作流程(5 个以上步骤、分支逻辑) |
| 工具/框架 | 卡夫卡、RabbitMQ、NATS、AWS EventBridge | Temporal.io、AWS Step Functions、Camunda、Axon |
隔离挑战与对策(处理“ACID minus I”)
由于 Saga 中的本地事务立即提交到其本地数据库,因此 Saga 模式缺乏与传统 ACID 保证的隔离 (I)。
如果客户端在 Saga 仍在运行时读取由 $T_1$ 修改的数据库行,则它们正在读取 未提交的中间状态。如果下游步骤失败并触发补偿 ($C_1$),则客户端已执行脏读。
缺乏隔离导致的常见异常
- 丢失更新:Saga A 更新一条记录。在 Saga A 完成之前,Saga B 会覆盖同一条记录。如果Saga A失败并执行补偿,它会覆盖Saga B的更新。
- 脏读:客户读取由 Saga A ($T_1$) 更新的可用库存。 Saga A 下游出现故障 ($T_3$) 并恢复库存 ($C_1$),但客户已根据过时数据下订单。
- 不可重复读取:服务在步骤 $T_1$ 读取数据,并在步骤 $T_3$ 再次读取数据,但另一个并发 Saga 修改了中间的数据。
对策和缓解策略
为了在缺乏隔离的情况下保持数据完整性,软件架构师实现了特定的隔离设计模式:
1. 语义锁(待处理/标记状态)
当本地事务 $T_1$ 更新数据库记录时,它将状态字段设置为 PENDING 或 APPROVAL_REQUIRED (例如 ORDER_PENDING_PAYMENT)。
读取此记录的其他并发 Sagas 必须检查语义锁定标志并阻止或更改其行为,直到状态更改为 COMMITTED 或 CANCELLED。
2. 提交订单
设计本地事务的顺序,以便在 Saga 执行的后期发生高风险或不可逆的操作,从而最大限度地减少漏洞窗口。
3.重读验证(乐观并发控制)
在执行关键步骤或补偿之前,重新读取目标数据库记录并验证版本时间戳(version_id)以确保不会发生并发修改。
4.悲观观点
重新排序 Saga 的步骤以最小化经济风险(例如,将支付授权尽可能靠近枢轴交易)。
生产序列流程:补偿事务
下面是完整的顺序流程图,说明了前向执行失败以及由此产生的后向补偿执行:
动手实践代码实现
让我们探索 Choreography(在 Go 中)和 Orchestration(在 Java Spring Boot 中)的生产就绪实现示例。
实现 1:Go 中基于编排的 Saga
在此 Go 示例中,我们演示了一个 订单服务 处理订单创建并通过事件代理侦听支付失败事件以执行补偿回滚逻辑。
package saga
import (
"context"
"encoding/json"
"fmt"
"log"
"time"
)
// Event definitions
type OrderCreatedEvent struct {
OrderID string `json:"order_id"`
CustomerID string `json:"customer_id"`
Amount float64 `json:"amount"`
}
type PaymentFailedEvent struct {
OrderID string `json:"order_id"`
Reason string `json:"reason"`
}
// OrderRepository handles local DB operations
type OrderRepository interface {
CreateOrder(ctx context.Context, orderID string, amount float64) error
UpdateOrderStatus(ctx context.Context, orderID string, status string) error
}
// EventBus abstraction for message broker (e.g., Kafka / NATS)
type EventBus interface {
Publish(topic string, payload []byte) error
Subscribe(topic string, handler func(payload []byte)) error
}
type OrderSagaChoreographer struct {
repo OrderRepository
eventBus EventBus
}
func NewOrderSagaChoreographer(repo OrderRepository, bus EventBus) *OrderSagaChoreographer {
c := &OrderSagaChoreographer{repo: repo, eventBus: bus}
c.registerSubscriptions()
return c
}
// Step 1: Forward Transaction (T1)
func (s *OrderSagaChoreographer) StartOrderSaga(ctx context.Context, orderID, customerID string, amount float64) error {
// Execute local database transaction
err := s.repo.CreateOrder(ctx, orderID, amount)
if err != nil {
return fmt.Errorf("failed local DB transaction T1: %w", err)
}
// Emit domain event for downstream Payment Service
event := OrderCreatedEvent{OrderID: orderID, CustomerID: customerID, Amount: amount}
bytes, _ := json.Marshal(event)
log.Printf("[SAGA][T1] Order %s created. Publishing OrderCreatedEvent...", orderID)
return s.eventBus.Publish("orders.created", bytes)
}
// Register subscription for compensating events
func (s *OrderSagaChoreographer) registerSubscriptions() {
_ = s.eventBus.Subscribe("payments.failed", func(payload []byte) {
var event PaymentFailedEvent
if err := json.Unmarshal(payload, &event); err != nil {
log.Printf("[ERROR] Corrupt payment event: %v", err)
return
}
// Execute Compensating Transaction (C1)
s.handlePaymentFailed(context.Background(), event)
})
}
// Step C1: Compensating Transaction
func (s *OrderSagaChoreographer) handlePaymentFailed(ctx context.Context, event PaymentFailedEvent) {
log.Printf("[SAGA][C1] Payment failed for Order %s (Reason: %s). Rolling back local order...", event.OrderID, event.Reason)
// Revert order status to CANCELLED in local DB
err := s.repo.UpdateOrderStatus(ctx, event.OrderID, "CANCELLED_PAYMENT_FAILED")
if err != nil {
log.Printf("[CRITICAL] Failed to execute compensating transaction C1 for Order %s: %v", event.OrderID, err)
// Trigger alert or write to Dead Letter Queue (DLQ)
return
}
log.Printf("[SAGA][SUCCESS] Order %s successfully compensated and cancelled.", event.OrderID)
}
实现2:Java中基于编排的Saga(Spring Boot)
在此 Java 示例中,我们使用状态机模式构建 Saga Orchestrator 来协调转发命令并在下游步骤失败时执行补偿回滚。
package com.ghaznix.saga.orchestrator;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;
import java.util.UUID;
public enum SagaState {
STARTED,
ORDER_CREATED,
PAYMENT_PROCESSED,
INVENTORY_RESERVED,
COMPLETED,
COMPENSATING_PAYMENT,
COMPENSATING_ORDER,
FAILED
}
@Service
public class OrderSagaOrchestrator {
private static final Logger log = LoggerFactory.getLogger(OrderSagaOrchestrator.class);
private final OrderServiceClient orderClient;
private final PaymentServiceClient paymentClient;
private final InventoryServiceClient inventoryClient;
public OrderSagaOrchestrator(OrderServiceClient orderClient,
PaymentServiceClient paymentClient,
InventoryServiceClient inventoryClient) {
this.orderClient = orderClient;
this.paymentClient = paymentClient;
this.inventoryClient = inventoryClient;
}
public boolean executeOrderSaga(String customerId, String productId, double amount, int quantity) {
String sagaId = UUID.randomUUID().toString();
log.info("[SAGA {}] Starting Order Saga Workflow...", sagaId);
SagaState currentState = SagaState.STARTED;
String orderId = null;
String paymentId = null;
try {
// Step 1: Forward Local Tx T1 - Create Order
log.info("[SAGA {}][T1] Sending CreateOrder command...", sagaId);
orderId = orderClient.createOrder(customerId, productId, amount);
currentState = SagaState.ORDER_CREATED;
// Step 2: Forward Local Tx T2 - Process Payment
log.info("[SAGA {}][T2] Sending ProcessPayment command for Order {}...", sagaId, orderId);
paymentId = paymentClient.chargePayment(orderId, customerId, amount);
currentState = SagaState.PAYMENT_PROCESSED;
// Step 3: Forward Local Tx T3 (Pivot Step) - Reserve Inventory
log.info("[SAGA {}][T3] Sending ReserveInventory command for Product {}...", sagaId, productId);
boolean stockReserved = inventoryClient.reserveStock(productId, quantity);
if (!stockReserved) {
throw new InventoryAllocationException("Stock allocation failed: Item out of stock.");
}
currentState = SagaState.INVENTORY_RESERVED;
log.info("[SAGA {}][SUCCESS] Saga completed successfully!", sagaId);
return true;
} catch (Exception ex) {
log.error("[SAGA {}][FAILURE] Step failed during state {}. Triggering Compensation...", sagaId, currentState, ex);
rollbackSaga(sagaId, currentState, orderId, paymentId);
return false;
}
}
private void rollbackSaga(String sagaId, SagaState failedState, String orderId, String paymentId) {
log.info("[SAGA {}] Initiating backward compensating transactions from state: {}", sagaId, failedState);
// Compensate Step 2 if payment was processed
if (failedState == SagaState.PAYMENT_PROCESSED || failedState == SagaState.INVENTORY_RESERVED) {
try {
log.info("[SAGA {}][C2] Executing Payment Refund compensation for Payment {}...", sagaId, paymentId);
paymentClient.refundPayment(paymentId);
} catch (Exception e) {
log.error("[CRITICAL][SAGA {}] Payment refund C2 failed! Manual intervention or DLQ required.", sagaId, e);
}
}
// Compensate Step 1 if order was created
if (failedState != SagaState.STARTED && orderId != null) {
try {
log.info("[SAGA {}][C1] Executing Order Cancel compensation for Order {}...", sagaId, orderId);
orderClient.cancelOrder(orderId);
} catch (Exception e) {
log.error("[CRITICAL][SAGA {}] Order cancellation C1 failed!", sagaId, e);
}
}
log.info("[SAGA {}] Saga compensation workflow finished. Final State: FAILED.", sagaId);
}
}
生产要点:将 Saga 与事务发件箱模式配对
在编排和编排中,执行本地事务 ($T_i$) 需要通过网络发布域事件或命令。
如果您的服务更新其 SQL 数据库,然后向 Kafka 发布消息,则数据库提交后的网络故障会导致静默事件丢失。相反,在数据库提交之前发布消息会导致幻象事件处理。
为了解决这个问题,Sagas 必须与事务性发件箱模式配对:
[Service Database Transaction Boundary]
+-----------------------------------------------------+
| 1. INSERT INTO business_table (orders/payments) |
| 2. INSERT INTO outbox_table (event_payload) |
+-----------------------------------------------------+
|
(CDC / Polling Message Relay)
|
v
[Message Broker / Kafka]
通过将事件有效负载保存到同一本地数据库事务内的 outbox 表中,可以保证原子性。异步后台进程(如 Debezium 或轮询中继)从发件箱表中读取数据并将事件可靠地发布到 Kafka。
此外,每个下游服务使用者必须实现幂等性(使用唯一的 idempotency_key 或消息重复数据删除标头),以便重试期间的重复消息传递不会触发重复费用或库存分配。
架构决策清单
在为微服务应用程序设计分布式事务时,请使用此实用决策矩阵:
Do you need cross-service data consistency?
|
+----------------+----------------+
| No | Yes
v v
Standard Single Service Can you accept Eventual
Local Database Consistency (BASE)?
|
+----------------+----------------+
| No | Yes
v v
Use Monolithic Core Adopt Saga Pattern
with Single ACID DB |
|
How complex is the workflow?
|
+-----------------+-----------------+
| Simple (2-3 steps) | Complex (4+ steps/branches)
v v
Choreography Saga Orchestration Saga
(Event-Driven Bus) (Temporal/Custom State Machine)
结论和要点
Saga 模式是一种重要的架构模式,用于管理跨微服务边界的分布式事务,而无需锁定资源或牺牲系统可用性。
摘要清单:
- 放弃云原生微服务中的2PC/XA:两阶段提交导致紧锁、高延迟和严重的可用性瓶颈。
- 将事务分解为本地步骤:将全局操作划分为本地事务 ($T_1 \dots T_n$) 与反向补偿事务 ($C_1 \dots C_{n-1}$) 配对。
- 选择正确的架构风格:
- 使用 Choreography 实现简单的、松散耦合的 2-3 步骤事件驱动流程。
- 对于需要集中可见性、分支和状态机跟踪的复杂业务工作流程使用编排。
- 实施隔离对策:通过使用语义锁(
PENDING标志)和重读乐观锁来防止脏读和丢失更新。 - 保证可靠的消息传递:始终将 Saga 实现与 事务发件箱模式 配对,并强制 幂等消费者 安全地处理重试。
通过精心实施 Saga 模式,您可以构建高度可用、可扩展的微服务,即使在网络发生故障时也能保持弹性和一致性。