返回博客

Saga 模式:分布式事务编排与补偿

在微服务架构中,一个业务操作可能跨越多个服务。传统的两阶段提交(2PC)在分布式环境下存在严重问题:

Coomia发布于 2025年12月18日16 分钟阅读
分享本文Twitter / X

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)在分布式环境下存在严重问题:

Code
场景:用户在 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):

Python
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 的执行:

Python
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 协调执行。以"创建风控规则"为例:

Python
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 提供了工作流状态的持久化,即使进程崩溃也能从中断点恢复:

Python
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 的执行进度、失败点和补偿状态:

Python
@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 表。

语义补偿:不是精确逆操作,而是业务语义上的"取消"。例如订单已发货,补偿操作不是"取消发货"(物理上已不可能),而是"创建退货单"。

Python
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 补偿操作的幂等性

补偿操作必须是幂等的——执行多次的效果和执行一次相同。这是因为在分布式环境中,补偿操作本身也可能失败并重试:

Python
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),等待人工介入:

Python
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 协调:

Code
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

大规模数据迁移需要分阶段执行,每个阶段都可以回滚:

Python
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 分配资源:

Code
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 具有完整的审计轨迹:

Python
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 的生命周期:

Code
CREATED → RUNNING → COMPLETED
                  ↘ COMPENSATING → COMPENSATED
                                 ↘ FAILED (需要人工介入)

#5.3 Saga + 幂等消费者

为了应对消息重复投递,Saga 的每个步骤都实现幂等性。结合幂等消费者模式,确保即使消息被重复处理,系统状态也保持一致。

#六、生产环境注意事项

#6.1 性能考量

Saga 模式相比单体事务有额外的性能开销:

开销来源影响缓解策略
多次网络调用延迟增加并行执行无依赖的步骤
事件持久化I/O 开销批量写入、异步持久化
补偿逻辑额外计算只在失败时触发
状态检查查询开销内存缓存活跃 Saga 状态

#6.2 并发控制

多个 Saga 实例可能同时操作同一资源。coomia-dip 使用分布式锁确保关键资源的互斥访问:

Python
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

  1. Saga 替代 2PC:在分布式微服务架构中,Saga 模式通过一系列可补偿的本地事务替代传统的两阶段提交
  2. 编排式优先:coomia-dip 默认使用编排式 Saga,通过 Temporal 工作流引擎实现持久化和可恢复性
  3. 补偿设计关键:补偿操作必须是幂等的、自包含的,不依赖外部状态假设
  4. 可观测性:每个 Saga 步骤都记录事件,提供完整的审计轨迹和实时监控
  5. 死信处理:无法自动补偿的 Saga 进入死信队列,需要人工介入
  6. 控制链长:避免过长的 Saga 链,超过 7 步考虑拆分为子 Saga

#Next Article

下一篇我们将深入探讨状态机模式(State Machine)——coomia-dip 如何用有限状态机管理 Ontology 对象的生命周期,以及状态机如何与 Saga 协调器配合工作。

S10-05: 状态机:对象生命周期管理

#Tags

#设计模式 #Saga #分布式事务 #补偿操作 #Temporal #编排式 #最终一致性 #跨Plane事务 #幂等性