Saga 模式:分布式事务编排与补偿
在微服务架构中,一个业务操作可能跨越多个服务。传统的两阶段提交(2PC)在分布式环境下存在严重问题:
Saga 模式:分布式事务编排与补偿
“系列:S10 设计模式 · 第 4 篇 | 难度:高级 | 阅读时间:18 分钟
#TL;DR
- Saga 模式将一个跨服务的长事务分解为一系列本地事务,每个本地事务都有对应的补偿操作。如果某一步失败,系统按反序执行补偿操作来撤销已完成的步骤。
- coomia-dip 在 Ontology Action 跨 Layer 执行时使用 Saga 模式保证最终一致性——从 Control Layer 的 Action 注册,到 Data Layer 的数据写入,再到 Intelligence Layer 的推理触发,任何环节失败都能安全回滚。
- 结合 Temporal 工作流引擎,coomia-dip 实现了可持久化、可恢复、可观测的 Saga 编排,支持编排式(Orchestration)和协作式(Choreography)两种模式。
#引言:分布式事务的困境
在微服务架构中,一个业务操作可能跨越多个服务。传统的两阶段提交(2PC)在分布式环境下存在严重问题:
场景:用户在 coomia-dip 中创建一个风控规则
1. Control Layer:注册 ObjectType 和 Action 定义
2. Data Layer:创建对应的 Iceberg 表结构
3. Intelligence Layer:初始化推理引擎中的规则模板
4. Agent Runtime:部署自动监控 Agent
如果第 3 步失败了怎么办?前两步已经提交,数据库已经有了记录。传统 2PC 在这种跨异构系统的场景下根本行不通——你无法要求 Kafka、Iceberg、Python 推理引擎都实现 XA 协议。
Saga 模式提供了优雅的解决方案:用一系列可补偿的本地事务替代全局事务。
#一、Saga 模式的核心概念
#1.1 本地事务与补偿操作
Saga 的核心思想是将全局事务拆分为一系列步骤,每个步骤包含一个前进操作(forward action)和一个补偿操作(compensating action):
from dataclasses import dataclass, field
from enum import Enum
from typing import Any
class SagaStepStatus(Enum):
PENDING = "pending"
RUNNING = "running"
COMPLETED = "completed"
COMPENSATING = "compensating"
COMPENSATED = "compensated"
FAILED = "failed"
@dataclass
class SagaStep:
"""A single step in a Saga."""
name: str
service: str # 目标服务/Layer
action: str # 前进操作
compensation: str # 补偿操作
input_data: dict[str, Any] = field(default_factory=dict)
output_data: dict[str, Any] = field(default_factory=dict)
status: SagaStepStatus = SagaStepStatus.PENDING
retry_count: int = 0
max_retries: int = 3
timeout_seconds: int = 30
@dataclass
class SagaDefinition:
"""Definition of a complete Saga."""
saga_id: str
name: str
steps: list[SagaStep]
compensation_policy: str = "backward" # backward | forward_recovery
timeout_seconds: int = 300
#1.2 Saga 执行协调器
coomia-dip 使用编排式(Orchestration)作为默认的 Saga 协调策略。一个中央协调器负责驱动整个 Saga 的执行:
class SagaOrchestrator:
"""Central Saga orchestrator for coomia-dip cross-Layer transactions."""
def __init__(
self,
event_store: "EventStore",
step_executors: dict[str, "StepExecutor"],
):
self._event_store = event_store
self._executors = step_executors
async def execute(self, saga: SagaDefinition) -> SagaResult:
"""Execute a Saga with automatic compensation on failure."""
completed_steps: list[SagaStep] = []
for step in saga.steps:
try:
step.status = SagaStepStatus.RUNNING
executor = self._executors[step.service]
result = await self._execute_with_retry(
executor, step, saga.timeout_seconds
)
step.output_data = result
step.status = SagaStepStatus.COMPLETED
completed_steps.append(step)
await self._event_store.append_saga_event(
saga.saga_id, step.name, "completed", result
)
except Exception as exc:
step.status = SagaStepStatus.FAILED
await self._event_store.append_saga_event(
saga.saga_id, step.name, "failed", {"error": str(exc)}
)
# 启动补偿
await self._compensate(saga, completed_steps)
return SagaResult(
saga_id=saga.saga_id,
status="compensated",
failed_step=step.name,
error=str(exc),
)
return SagaResult(saga_id=saga.saga_id, status="completed")
async def _compensate(
self, saga: SagaDefinition, completed_steps: list[SagaStep]
) -> None:
"""Execute compensation in reverse order."""
for step in reversed(completed_steps):
try:
step.status = SagaStepStatus.COMPENSATING
executor = self._executors[step.service]
await executor.compensate(step.compensation, step.output_data)
step.status = SagaStepStatus.COMPENSATED
await self._event_store.append_saga_event(
saga.saga_id, step.name, "compensated", {}
)
except Exception as comp_exc:
step.status = SagaStepStatus.FAILED
await self._event_store.append_saga_event(
saga.saga_id,
step.name,
"compensation_failed",
{"error": str(comp_exc)},
)
# 补偿失败需要人工介入
await self._alert_manual_intervention(saga, step, comp_exc)
async def _execute_with_retry(
self, executor: "StepExecutor", step: SagaStep, timeout: int
) -> dict[str, Any]:
"""Execute a step with retry logic."""
import asyncio
last_error: Exception | None = None
for attempt in range(step.max_retries + 1):
try:
return await asyncio.wait_for(
executor.execute(step.action, step.input_data),
timeout=min(step.timeout_seconds, timeout),
)
except asyncio.TimeoutError:
last_error = TimeoutError(
f"Step {step.name} timed out after {step.timeout_seconds}s"
)
except Exception as e:
last_error = e
if attempt < step.max_retries:
await asyncio.sleep(2 ** attempt) # 指数退避
raise last_error # type: ignore[misc]
async def _alert_manual_intervention(
self, saga: SagaDefinition, step: SagaStep, error: Exception
) -> None:
"""Alert for manual intervention when compensation fails."""
# 发送告警通知运维团队
pass
#1.3 编排式 vs 协作式
Saga 有两种协调模式,coomia-dip 支持两种但默认使用编排式:
| 维度 | 编排式(Orchestration) | 协作式(Choreography) |
|---|---|---|
| 协调方式 | 中央协调器驱动 | 事件驱动,服务自行监听 |
| 耦合度 | 协调器知道所有步骤 | 服务只知道自己的事件 |
| 可观测性 | 高——协调器记录全局状态 | 低——需要聚合多源事件 |
| 复杂度 | 流程集中管理,易理解 | 分散在各服务,难追踪 |
| 适用场景 | 跨 Layer 事务、复杂流程 | 简单通知链、松耦合集成 |
| coomia-dip 使用 | Action 执行、规则部署 | 事件通知、监控触发 |
#二、coomia-dip 中的 Saga 实践
#2.1 Action 跨 Layer 执行 Saga
当用户通过 SDK 调用一个 Ontology Action,这个 Action 可能需要跨多个 Layer 协调执行。以"创建风控规则"为例:
class CreateRiskRuleSaga:
"""Saga for creating a risk control rule across Layers."""
def build(self, rule_config: dict) -> SagaDefinition:
return SagaDefinition(
saga_id=generate_id(),
name="create_risk_rule",
steps=[
SagaStep(
name="register_object_type",
service="control-Layer",
action="register_risk_rule_type",
compensation="unregister_risk_rule_type",
input_data={"schema": rule_config["schema"]},
),
SagaStep(
name="create_storage",
service="data-Layer",
action="create_iceberg_table",
compensation="drop_iceberg_table",
input_data={"table_spec": rule_config["storage"]},
),
SagaStep(
name="init_reasoning_template",
service="intelligence-Layer",
action="deploy_rule_template",
compensation="undeploy_rule_template",
input_data={"template": rule_config["reasoning"]},
),
SagaStep(
name="setup_monitoring",
service="agent-runtime",
action="create_monitor_agent",
compensation="destroy_monitor_agent",
input_data={"monitor": rule_config["monitoring"]},
),
],
timeout_seconds=120,
)
#2.2 基于 Temporal 的 Saga 持久化
coomia-dip 使用 Temporal 作为 Saga 的持久化和恢复引擎。Temporal 提供了工作流状态的持久化,即使进程崩溃也能从中断点恢复:
from temporalio import workflow, activity
from datetime import timedelta
@activity.defn
async def register_object_type(input_data: dict) -> dict:
"""Register ObjectType in Control Layer via gRPC."""
async with grpc_channel("control-Layer:50051") as channel:
stub = OntologyServiceStub(channel)
response = await stub.RegisterObjectType(
RegisterObjectTypeRequest(**input_data)
)
return {"type_id": response.type_id, "version": response.version}
@activity.defn
async def unregister_object_type(output_data: dict) -> None:
"""Compensate: unregister the ObjectType."""
async with grpc_channel("control-Layer:50051") as channel:
stub = OntologyServiceStub(channel)
await stub.UnregisterObjectType(
UnregisterObjectTypeRequest(type_id=output_data["type_id"])
)
@activity.defn
async def create_iceberg_table(input_data: dict) -> dict:
"""Create Iceberg table in Data Layer."""
async with grpc_channel("data-Layer:50052") as channel:
stub = DataServiceStub(channel)
response = await stub.CreateTable(
CreateTableRequest(**input_data)
)
return {"table_id": response.table_id, "location": response.location}
@activity.defn
async def drop_iceberg_table(output_data: dict) -> None:
"""Compensate: drop the Iceberg table."""
async with grpc_channel("data-Layer:50052") as channel:
stub = DataServiceStub(channel)
await stub.DropTable(
DropTableRequest(table_id=output_data["table_id"])
)
@workflow.defn
class CreateRiskRuleWorkflow:
"""Temporal workflow implementing the Create Risk Rule Saga."""
@workflow.run
async def run(self, rule_config: dict) -> dict:
completed: list[dict] = []
try:
# Step 1: Register ObjectType
type_result = await workflow.execute_activity(
register_object_type,
{"schema": rule_config["schema"]},
start_to_close_timeout=timedelta(seconds=30),
retry_policy=RetryPolicy(maximum_attempts=3),
)
completed.append(("unregister_object_type", type_result))
# Step 2: Create storage
storage_result = await workflow.execute_activity(
create_iceberg_table,
{"table_spec": rule_config["storage"]},
start_to_close_timeout=timedelta(seconds=30),
)
completed.append(("drop_iceberg_table", storage_result))
# Step 3: Deploy reasoning template
reasoning_result = await workflow.execute_activity(
deploy_rule_template,
{"template": rule_config["reasoning"]},
start_to_close_timeout=timedelta(seconds=60),
)
completed.append(("undeploy_rule_template", reasoning_result))
# Step 4: Setup monitoring
monitor_result = await workflow.execute_activity(
create_monitor_agent,
{"monitor": rule_config["monitoring"]},
start_to_close_timeout=timedelta(seconds=30),
)
return {
"status": "completed",
"type_id": type_result["type_id"],
"table_id": storage_result["table_id"],
}
except Exception as e:
# 补偿已完成的步骤
for comp_action, comp_data in reversed(completed):
try:
await workflow.execute_activity(
comp_action,
comp_data,
start_to_close_timeout=timedelta(seconds=30),
)
except Exception as comp_error:
workflow.logger.error(
f"Compensation failed: {comp_action}, "
f"error: {comp_error}"
)
raise
#2.3 Saga 状态可视化
coomia-dip 的 Platform Console 提供 Saga 执行的实时可视化,运维人员可以看到每个 Saga 的执行进度、失败点和补偿状态:
@dataclass
class SagaVisualization:
"""Saga execution visualization data."""
saga_id: str
name: str
status: str
started_at: datetime
completed_at: datetime | None
steps: list[StepVisualization]
total_duration_ms: int
compensation_duration_ms: int | None
@dataclass
class StepVisualization:
"""Individual step visualization."""
name: str
service: str
status: str
started_at: datetime
completed_at: datetime | None
duration_ms: int
retry_count: int
error: str | None
class SagaDashboardService:
"""Service for Saga monitoring dashboard."""
async def get_saga_history(
self,
tenant_id: str,
time_range: tuple[datetime, datetime],
status_filter: list[str] | None = None,
) -> list[SagaVisualization]:
"""Retrieve Saga execution history with visualization data."""
query = (
f"SELECT * FROM saga_executions "
f"WHERE tenant_id = '{tenant_id}' "
f"AND started_at BETWEEN '{time_range[0]}' AND '{time_range[1]}'"
)
if status_filter:
statuses = "', '".join(status_filter)
query += f" AND status IN ('{statuses}')"
# ... 查询执行
pass
async def get_compensation_stats(
self, tenant_id: str
) -> dict[str, Any]:
"""Get compensation statistics for monitoring."""
return {
"total_sagas": 0,
"completed": 0,
"compensated": 0,
"compensation_rate": 0.0,
"avg_compensation_time_ms": 0,
"manual_interventions": 0,
}
#三、补偿策略设计
#3.1 语义补偿 vs 精确回滚
在分布式系统中,补偿操作不总是精确的"撤销"。coomia-dip 区分两种补偿策略:
精确回滚:操作的逆操作,完全恢复到操作前状态。例如删除刚创建的 Iceberg 表。
语义补偿:不是精确逆操作,而是业务语义上的"取消"。例如订单已发货,补偿操作不是"取消发货"(物理上已不可能),而是"创建退货单"。
class CompensationStrategies:
"""Different compensation strategies for various scenarios."""
@staticmethod
async def exact_rollback(step_name: str, output_data: dict) -> None:
"""Exact inverse operation — fully restores prior state."""
# 例如:删除刚创建的资源
pass
@staticmethod
async def semantic_compensation(step_name: str, output_data: dict) -> None:
"""Business-level compensation — not exact inverse."""
# 例如:标记资源为"已取消"而非删除
pass
@staticmethod
async def idempotent_retry(step_name: str, input_data: dict) -> None:
"""Retry with idempotency key instead of compensating."""
# 对于幂等操作,重试可能比补偿更合适
pass
#3.2 补偿操作的幂等性
补偿操作必须是幂等的——执行多次的效果和执行一次相同。这是因为在分布式环境中,补偿操作本身也可能失败并重试:
class IdempotentCompensation:
"""Ensure compensation operations are idempotent."""
def __init__(self, state_store: "StateStore"):
self._state_store = state_store
async def compensate_with_idempotency(
self,
saga_id: str,
step_name: str,
compensation_fn: callable,
data: dict,
) -> None:
"""Execute compensation with idempotency guarantee."""
idempotency_key = f"{saga_id}:{step_name}:compensation"
# 检查是否已经执行过
if await self._state_store.is_completed(idempotency_key):
return # 已执行,跳过
try:
await compensation_fn(data)
await self._state_store.mark_completed(idempotency_key)
except Exception:
# 记录失败,等待重试
await self._state_store.mark_failed(idempotency_key)
raise
#3.3 超时与死信处理
当补偿操作反复失败时,coomia-dip 将 Saga 实例移入死信队列(Dead Letter Queue),等待人工介入:
class SagaDeadLetterHandler:
"""Handle Sagas that cannot be automatically compensated."""
async def move_to_dead_letter(
self, saga_id: str, failed_step: str, error: str
) -> None:
"""Move a failed Saga to the dead letter queue."""
await self._dead_letter_store.save({
"saga_id": saga_id,
"failed_step": failed_step,
"error": error,
"timestamp": datetime.utcnow(),
"requires_manual_intervention": True,
})
# 通知运维
await self._notification_service.send_alert(
severity="critical",
title=f"Saga compensation failed: {saga_id}",
message=(
f"Step '{failed_step}' compensation failed after all retries. "
f"Manual intervention required. Error: {error}"
),
)
async def retry_dead_letter(self, saga_id: str) -> SagaResult:
"""Manual retry of a dead-lettered Saga."""
saga_state = await self._dead_letter_store.get(saga_id)
# 从失败点重新尝试补偿
return await self._orchestrator.resume_compensation(
saga_id, saga_state["failed_step"]
)
#四、Saga 在 coomia-dip 各场景中的应用
#4.1 Ontology Schema 变更 Saga
当 Ontology 的 Schema 发生变更(如新增属性、修改关系),需要跨 Layer 协调:
Step 1: Control Layer — 验证 Schema 兼容性
Step 2: Control Layer — 更新 ObjectType 定义
Step 3: Data Layer — 执行 Iceberg Schema Evolution
Step 4: Intelligence Layer — 更新推理引擎的属性映射
Step 5: SDK Layer — 重新生成 SDK 类型定义
补偿:
Step 5 失败 → 回滚 SDK 代码生成
Step 4 失败 → 恢复推理引擎旧映射
Step 3 失败 → 回滚 Iceberg Schema
Step 2 失败 → 恢复 ObjectType 旧定义
#4.2 数据迁移 Saga
大规模数据迁移需要分阶段执行,每个阶段都可以回滚:
class DataMigrationSaga:
"""Saga for large-scale data migration."""
def build(self, migration_plan: dict) -> SagaDefinition:
return SagaDefinition(
saga_id=generate_id(),
name="data_migration",
steps=[
SagaStep(
name="validate_source",
service="data-Layer",
action="validate_source_data",
compensation="noop", # 验证无需补偿
timeout_seconds=60,
),
SagaStep(
name="create_target_schema",
service="data-Layer",
action="create_migration_target",
compensation="drop_migration_target",
timeout_seconds=30,
),
SagaStep(
name="copy_data",
service="data-Layer",
action="copy_data_batch",
compensation="delete_copied_data",
timeout_seconds=3600, # 大数据量可能需要较长时间
),
SagaStep(
name="validate_target",
service="data-Layer",
action="validate_target_data",
compensation="noop",
timeout_seconds=120,
),
SagaStep(
name="switch_references",
service="control-Layer",
action="update_ontology_references",
compensation="revert_ontology_references",
timeout_seconds=30,
),
],
timeout_seconds=7200,
)
#4.3 多租户资源分配 Saga
新租户开通时需要在多个 Layer 分配资源:
Step 1: Control Layer — 创建租户元数据
Step 2: Data Layer — 分配存储命名空间(Nessie Branch)
Step 3: Intelligence Layer — 初始化推理引擎实例
Step 4: Agent Runtime — 部署默认 Agent
Step 5: SDK Layer — 生成租户专属 API 密钥
任何步骤失败 → 按反序清理已分配的资源
#五、Saga 与其他模式的结合
#5.1 Saga + Event Sourcing
Saga 的每一步执行和补偿都记录为事件,存入 Event Store。这使得 Saga 具有完整的审计轨迹:
class SagaEventStore:
"""Record Saga events for audit trail."""
async def append_saga_event(
self, saga_id: str, step_name: str, action: str, data: dict
) -> None:
event = DomainEvent(
event_id=generate_id(),
event_type=f"saga.{action}",
aggregate_id=saga_id,
aggregate_type="saga",
sequence_number=await self._next_sequence(saga_id),
timestamp=datetime.utcnow(),
payload={
"step_name": step_name,
"action": action,
"data": data,
},
metadata=EventMetadata(
actor_id="system",
actor_type="saga_orchestrator",
tenant_id=self._current_tenant_id,
world_id=self._current_world_id,
source_plane="orchestration",
trace_id=self._current_trace_id,
),
)
await self._event_store.append([event])
#5.2 Saga + State Machine
每个 Saga 实例本质上是一个状态机。coomia-dip 使用状态机(下一篇将详细介绍)来管理 Saga 的生命周期:
CREATED → RUNNING → COMPLETED
↘ COMPENSATING → COMPENSATED
↘ FAILED (需要人工介入)
#5.3 Saga + 幂等消费者
为了应对消息重复投递,Saga 的每个步骤都实现幂等性。结合幂等消费者模式,确保即使消息被重复处理,系统状态也保持一致。
#六、生产环境注意事项
#6.1 性能考量
Saga 模式相比单体事务有额外的性能开销:
| 开销来源 | 影响 | 缓解策略 |
|---|---|---|
| 多次网络调用 | 延迟增加 | 并行执行无依赖的步骤 |
| 事件持久化 | I/O 开销 | 批量写入、异步持久化 |
| 补偿逻辑 | 额外计算 | 只在失败时触发 |
| 状态检查 | 查询开销 | 内存缓存活跃 Saga 状态 |
#6.2 并发控制
多个 Saga 实例可能同时操作同一资源。coomia-dip 使用分布式锁确保关键资源的互斥访问:
class SagaResourceLock:
"""Distributed lock for Saga resource protection."""
async def acquire(
self, resource_id: str, saga_id: str, timeout: int = 30
) -> bool:
"""Acquire a distributed lock for a resource."""
return await self._redis.set(
f"saga:lock:{resource_id}",
saga_id,
nx=True,
ex=timeout,
)
async def release(self, resource_id: str, saga_id: str) -> None:
"""Release a distributed lock."""
current = await self._redis.get(f"saga:lock:{resource_id}")
if current == saga_id:
await self._redis.delete(f"saga:lock:{resource_id}")
#6.3 监控与告警
生产环境中必须监控 Saga 的关键指标:
- 补偿率:补偿执行次数 / 总 Saga 次数。过高说明下游服务不稳定。
- 平均执行时间:Saga 从开始到完成的平均耗时。
- 死信队列深度:堆积的失败 Saga 数量。
- 步骤失败分布:哪些步骤最常失败。
#七、反模式与注意事项
#7.1 不要在 Saga 中使用隔离性假设
Saga 不提供隔离性(Isolation)。在 Saga 执行期间,中间状态对其他事务可见。设计时需要考虑:
- 脏读:其他事务可能读到 Saga 执行中的中间状态
- 丢失更新:并发 Saga 可能覆盖彼此的更新
- 解决方案:语义锁(Semantic Lock)、交换律补偿(Commutative Updates)
#7.2 避免过长的 Saga 链
步骤越多,失败概率越高,补偿越复杂。经验法则:
- 超过 7 个步骤的 Saga 应考虑拆分
- 每个步骤的超时应合理设置
- 可以用嵌套 Saga(子 Saga)管理复杂流程
#7.3 补偿操作不能依赖外部状态变化
补偿操作应该只依赖于对应前进操作的输出数据,不应该假设外部环境没有变化。补偿发生时,外部状态可能已经因为其他操作而改变。
#Key Takeaways
- Saga 替代 2PC:在分布式微服务架构中,Saga 模式通过一系列可补偿的本地事务替代传统的两阶段提交
- 编排式优先:coomia-dip 默认使用编排式 Saga,通过 Temporal 工作流引擎实现持久化和可恢复性
- 补偿设计关键:补偿操作必须是幂等的、自包含的,不依赖外部状态假设
- 可观测性:每个 Saga 步骤都记录事件,提供完整的审计轨迹和实时监控
- 死信处理:无法自动补偿的 Saga 进入死信队列,需要人工介入
- 控制链长:避免过长的 Saga 链,超过 7 步考虑拆分为子 Saga
#Next Article
下一篇我们将深入探讨状态机模式(State Machine)——coomia-dip 如何用有限状态机管理 Ontology 对象的生命周期,以及状态机如何与 Saga 协调器配合工作。
S10-05: 状态机:对象生命周期管理
#Tags
#设计模式 #Saga #分布式事务 #补偿操作 #Temporal #编排式 #最终一致性 #跨Plane事务 #幂等性