返回博客

Temporal 工作流引擎深潜(Part 1):持久化执行、Activity 重试与 Saga 补偿

1. [Temporal 在 coomia-dip 中的定位](#1-temporal-在-coomia-dip-中的定位)

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

系列:S8 技术组件深潜 · 第 8 篇 | 难度:高级 | 阅读时间:20 分钟

Temporal 工作流引擎深潜(Part 1):持久化执行、Activity 重试与 Saga 补偿

#TL;DR

  • Temporal 是 coomia-dip Agent Runtime(Agent Runtime Layer)的工作流编排核心,为长时间运行的决策流程提供持久化执行保证
  • 本文深入分析 Temporal Server 的架构(Frontend/History/Matching/Worker)、Event Sourcing 模型、Activity 重试策略的数学模型、以及 Saga 补偿模式在 Action 执行链中的应用
  • 涵盖 Workflow Determinism 约束、Signal/Query 通信模式、Child Workflow 编排模式,以及 coomia-dip 中 6 种典型工作流模板

#目录

  1. Temporal 在 coomia-dip 中的定位
  2. Temporal Server 四组件架构
  3. Event Sourcing 持久化模型
  4. Workflow Determinism 约束
  5. Activity 重试策略的数学模型
  6. Saga 补偿模式
  7. Signal 与 Query 通信模式
  8. Child Workflow 编排模式
  9. coomia-dip 6 种工作流模板
  10. 生产部署与调优
  11. Key Takeaways

#1. Temporal 在 coomia-dip 中的定位

#1.1 为什么选择 Temporal?

在 coomia-dip 的 Agent Runtime(Agent Runtime Layer)中,存在大量需要持久化执行保证的场景:

场景持续时间失败恢复需求补偿需求
Action 审批链分钟-天
数据管道编排小时
Agent 多步推理秒-分钟
批量数据导入小时
定时报告生成分钟
跨系统数据同步持续运行极高

传统方案的局限:

  • 消息队列 + 状态机:需要手动管理状态持久化、重试、补偿
  • Cron + 数据库锁:无法处理长时间运行的流程
  • DolphinScheduler:适合批处理调度,不适合事件驱动的实时工作流

Temporal 提供:

  • 自动持久化执行状态(即使进程崩溃也能恢复)
  • 声明式重试策略(指数退避、最大重试次数)
  • 原生 Saga 补偿支持
  • 毫秒级 Timer 精度
  • 可视化工作流执行历史

#1.2 coomia-dip 的 Temporal 架构

Code
┌─────────────────────────────────────────────────┐
│                  Agent Runtime Layer: Agent Runtime          │
│  ┌──────────┐  ┌──────────┐  ┌──────────┐      │
│  │ Action   │  │ Pipeline │  │ Agent    │      │
│  │ Workflow │  │ Workflow │  │ Workflow │      │
│  └────┬─────┘  └────┬─────┘  └────┬─────┘      │
│       │              │              │            │
│  ┌────┴──────────────┴──────────────┴────┐      │
│  │         Temporal Python SDK            │      │
│  │    (Worker / Client / Activities)      │      │
│  └────────────────┬──────────────────────┘      │
└───────────────────┼──────────────────────────────┘
                    │ gRPC
┌───────────────────┼──────────────────────────────┐
│  Temporal Server  │                              │
│  ┌─────────┐  ┌──┴──────┐  ┌──────────┐        │
│  │Frontend │  │ History │  │ Matching │        │
│  │ Service │  │ Service │  │ Service  │        │
│  └─────────┘  └─────────┘  └──────────┘        │
│                     │                            │
│  ┌──────────────────┴───────────────────┐       │
│  │    Persistence (PostgreSQL)           │       │
│  └──────────────────────────────────────┘       │
└──────────────────────────────────────────────────┘

#2. Temporal Server 四组件架构

#2.1 Frontend Service

Frontend Service 是所有 gRPC 请求的入口点:

Code
Client SDK ──gRPC──→ Frontend Service
                        │
                        ├── Rate Limiting (令牌桶)
                        ├── Request Validation
                        ├── Namespace Routing
                        └── 转发到 History / Matching Service

关键职责

  • 速率限制:基于 Namespace 的令牌桶限流,防止单租户过载
  • 请求校验:Workflow ID 格式、Payload 大小限制(默认 2MB)
  • Namespace 隔离:coomia-dip 为每个 World 创建独立 Namespace
Python
# coomia-dip Namespace 管理
from temporalio.client import Client

async def create_world_namespace(world_id: str) -> None:
    client = await Client.connect("temporal-server:7233")
    # 每个 World 使用独立的 Namespace 实现隔离
    # Namespace 命名规则:coomia-dip-{world_id}
    await client.operator_service.create_namespace(
        name=f"coomia-dip-{world_id}",
        retention_period=timedelta(days=30),  # 工作流历史保留 30 天
    )

#2.2 History Service

History Service 是 Temporal 的核心组件,负责工作流执行状态管理:

Code
Workflow Execution
    │
    ├── Shard 1 ──→ History Service Instance A
    ├── Shard 2 ──→ History Service Instance B
    ├── Shard 3 ──→ History Service Instance A  (多 Shard 可映射到同实例)
    └── Shard N ──→ History Service Instance C

分片策略

  • 默认 512 个 Shard,按 Workflow ID 哈希分配
  • 每个 Shard 由一个 History Service 实例独占处理(避免并发冲突)
  • Shard 所有权通过 Lease 机制在实例间转移

事件持久化

SQL
-- History Event 存储结构
CREATE TABLE executions_v2 (
    shard_id        INT NOT NULL,
    namespace_id    BINARY(16) NOT NULL,
    workflow_id     VARCHAR(255) NOT NULL,
    run_id          BINARY(16) NOT NULL,
    event_id        BIGINT NOT NULL,
    event_type      INT NOT NULL,
    event_payload   BLOB NOT NULL,
    PRIMARY KEY (shard_id, namespace_id, workflow_id, run_id, event_id)
);

#2.3 Matching Service

Matching Service 实现 Task Queue 的分发机制:

Code
Worker Poll ──→ Matching Service ──→ 返回 Workflow/Activity Task

Workflow Task Queue:
  ┌─────────────────────────────────────┐
  │ Task 1: WF-abc, Decision needed     │ → Worker A polls
  │ Task 2: WF-def, Decision needed     │ → Worker B polls
  │ Task 3: WF-ghi, Decision needed     │ → Worker A polls (round-robin)
  └─────────────────────────────────────┘

Activity Task Queue:
  ┌─────────────────────────────────────┐
  │ Task 1: Execute ApprovalActivity    │ → Worker C polls
  │ Task 2: Execute NotifyActivity      │ → Worker D polls
  └─────────────────────────────────────┘

Task Queue 类型

  • Sync Match:Worker 已在等待,Task 直接分发(延迟最低)
  • Async Match:Task 先写入数据库,Worker 后续 Poll 获取

#2.4 Worker 进程

coomia-dip 的 Temporal Worker 使用 Python SDK:

Python
import asyncio
from temporalio.client import Client
from temporalio.worker import Worker

from coomia-dip.workflows.action_workflow import ActionApprovalWorkflow
from coomia-dip.activities.action_activities import (
    validate_action,
    execute_action,
    notify_approvers,
    record_audit_log,
)

async def run_worker():
    client = await Client.connect("temporal-server:7233")

    worker = Worker(
        client,
        task_queue="coomia-dip-action-queue",
        workflows=[ActionApprovalWorkflow],
        activities=[
            validate_action,
            execute_action,
            notify_approvers,
            record_audit_log,
        ],
        max_concurrent_workflow_tasks=100,
        max_concurrent_activities=50,
        max_cached_workflows=500,
    )

    await worker.run()

if __name__ == "__main__":
    asyncio.run(run_worker())

#3. Event Sourcing 持久化模型

#3.1 Workflow 执行即事件序列

Temporal 将每个 Workflow 的执行记录为不可变的事件序列:

Code
Event History for Workflow "action-approval-123":

  Event 1:  WorkflowExecutionStarted
  Event 2:  WorkflowTaskScheduled
  Event 3:  WorkflowTaskStarted
  Event 4:  WorkflowTaskCompleted
  Event 5:  ActivityTaskScheduled (validate_action)
  Event 6:  ActivityTaskStarted
  Event 7:  ActivityTaskCompleted (result: valid)
  Event 8:  WorkflowTaskScheduled
  Event 9:  WorkflowTaskStarted
  Event 10: WorkflowTaskCompleted
  Event 11: TimerStarted (wait for approval, 24h)
  Event 12: SignalExternalWorkflowExecutionInitiated
  ...
  Event 25: WorkflowExecutionCompleted

#3.2 Replay 恢复机制

当 Worker 崩溃重启后,Workflow 通过 Replay Event History 恢复状态:

Python
# Workflow 代码(所有决策逻辑在 Replay 时重新执行)
@workflow.defn
class ActionApprovalWorkflow:
    @workflow.run
    async def run(self, request: ActionRequest) -> ActionResult:
        # Step 1: 验证 — Replay 时跳过 Activity 执行,直接使用记录的结果
        validation = await workflow.execute_activity(
            validate_action,
            request,
            start_to_close_timeout=timedelta(seconds=30),
        )

        # Step 2: 通知审批人
        await workflow.execute_activity(
            notify_approvers,
            request.approvers,
            start_to_close_timeout=timedelta(seconds=10),
        )

        # Step 3: 等待审批(可能等待数天)
        # Replay 时如果已有 Signal Event,直接返回结果
        approval = await workflow.wait_condition(
            lambda: self._approval_decision is not None,
            timeout=timedelta(hours=24),
        )

        if self._approval_decision == "approved":
            result = await workflow.execute_activity(
                execute_action,
                request,
                start_to_close_timeout=timedelta(minutes=5),
            )
            return ActionResult(status="completed", data=result)
        else:
            return ActionResult(status="rejected")

Replay 的关键规则

  1. Workflow 代码的决策路径在 Replay 时必须产生相同结果(Determinism)
  2. Activity 的实际执行不会重复 — 使用 Event History 中记录的结果
  3. Timer 在 Replay 时立即完成(如果已过期)

#3.3 Event History 大小管理

Python
# 大 Event History 的优化:使用 Continue-As-New
@workflow.defn
class LongRunningPipelineWorkflow:
    @workflow.run
    async def run(self, state: PipelineState) -> PipelineResult:
        while not state.is_complete:
            batch = await workflow.execute_activity(
                process_next_batch,
                state,
                start_to_close_timeout=timedelta(minutes=10),
            )
            state.update(batch)
            state.iterations += 1

            # 每 1000 次迭代,使用 Continue-As-New 重置 Event History
            if state.iterations % 1000 == 0:
                workflow.continue_as_new(state)

        return PipelineResult(state)

Event History 限制

指标默认限制coomia-dip 配置
最大事件数50,00050,000
最大 History 大小50 MB50 MB
建议 Continue-As-New 阈值10,000 事件5,000 事件

#4. Workflow Determinism 约束

#4.1 什么是 Determinism?

Workflow 代码在 Replay 时必须产生与首次执行完全相同的决策序列。以下操作违反 Determinism:

Python
# ❌ 违反 Determinism 的代码

@workflow.defn
class BadWorkflow:
    @workflow.run
    async def run(self, request: dict) -> str:
        # ❌ 使用系统时间(Replay 时时间不同)
        if datetime.now() > some_deadline:
            pass

        # ❌ 使用随机数(Replay 时结果不同)
        if random.random() > 0.5:
            pass

        # ❌ 使用 UUID(Replay 时生成不同值)
        task_id = str(uuid.uuid4())

        # ❌ 直接 I/O 操作(文件/网络/数据库)
        data = requests.get("http://api.example.com/data")

        # ❌ 使用可变全局状态
        global_counter += 1
Python
# ✅ 正确的 Deterministic 代码

@workflow.defn
class GoodWorkflow:
    @workflow.run
    async def run(self, request: dict) -> str:
        # ✅ 使用 Temporal 提供的时间
        now = workflow.now()

        # ✅ 使用 Temporal 提供的随机数(Replay 安全)
        value = workflow.random().random()

        # ✅ 使用 Temporal 的 Side Effect 生成 UUID
        task_id = await workflow.execute_activity(
            generate_task_id,
            start_to_close_timeout=timedelta(seconds=5),
        )

        # ✅ I/O 操作放在 Activity 中
        data = await workflow.execute_activity(
            fetch_data,
            "http://api.example.com/data",
            start_to_close_timeout=timedelta(seconds=30),
        )

#4.2 版本化(Patching)

当需要修改已运行 Workflow 的逻辑时:

Python
@workflow.defn
class EvolvingWorkflow:
    @workflow.run
    async def run(self, request: ActionRequest) -> ActionResult:
        # 版本化:旧 Workflow 走旧路径,新 Workflow 走新路径
        if workflow.patched("add-risk-check"):
            # 新版本:增加风险检查步骤
            risk = await workflow.execute_activity(
                check_risk_level,
                request,
                start_to_close_timeout=timedelta(seconds=30),
            )
            if risk.level == "HIGH":
                request.require_extra_approval = True

        # 后续逻辑不变
        validation = await workflow.execute_activity(
            validate_action,
            request,
            start_to_close_timeout=timedelta(seconds=30),
        )

#5. Activity 重试策略的数学模型

#5.1 重试配置

Python
from temporalio.common import RetryPolicy

# coomia-dip 标准重试策略
STANDARD_RETRY = RetryPolicy(
    initial_interval=timedelta(seconds=1),      # 首次重试间隔
    backoff_coefficient=2.0,                      # 指数退避系数
    maximum_interval=timedelta(minutes=5),        # 最大重试间隔
    maximum_attempts=10,                           # 最大重试次数
    non_retryable_error_types=[                   # 不重试的错误类型
        "ValidationError",
        "PermissionDeniedError",
        "NotFoundError",
    ],
)

# 使用重试策略
result = await workflow.execute_activity(
    execute_action,
    request,
    start_to_close_timeout=timedelta(minutes=5),
    retry_policy=STANDARD_RETRY,
)

#5.2 退避时间计算

第 N 次重试的等待时间:

Code
wait(N) = min(initial_interval × backoff_coefficient^(N-1), maximum_interval)

示例(initial=1s, coefficient=2.0, max=300s):
  重试 1: min(1 × 2^0, 300) = 1s
  重试 2: min(1 × 2^1, 300) = 2s
  重试 3: min(1 × 2^2, 300) = 4s
  重试 4: min(1 × 2^3, 300) = 8s
  重试 5: min(1 × 2^4, 300) = 16s
  重试 6: min(1 × 2^5, 300) = 32s
  重试 7: min(1 × 2^6, 300) = 64s
  重试 8: min(1 × 2^7, 300) = 128s
  重试 9: min(1 × 2^8, 300) = 256s
  重试 10: min(1 × 2^9, 300) = 300s (capped)

总等待时间 ≈ 811s ≈ 13.5 分钟

#5.3 不同场景的重试策略

Python
# 场景 1:调用外部 API(可能暂时不可用)
EXTERNAL_API_RETRY = RetryPolicy(
    initial_interval=timedelta(seconds=2),
    backoff_coefficient=3.0,
    maximum_interval=timedelta(minutes=10),
    maximum_attempts=15,
    non_retryable_error_types=["AuthenticationError"],
)

# 场景 2:数据库操作(短暂连接问题)
DB_RETRY = RetryPolicy(
    initial_interval=timedelta(milliseconds=100),
    backoff_coefficient=2.0,
    maximum_interval=timedelta(seconds=30),
    maximum_attempts=5,
)

# 场景 3:幂等写入(可以安全重试)
IDEMPOTENT_WRITE_RETRY = RetryPolicy(
    initial_interval=timedelta(seconds=1),
    backoff_coefficient=2.0,
    maximum_interval=timedelta(minutes=1),
    maximum_attempts=20,
)

# 场景 4:非幂等操作(需要谨慎)
NON_IDEMPOTENT_RETRY = RetryPolicy(
    initial_interval=timedelta(seconds=5),
    backoff_coefficient=1.5,
    maximum_interval=timedelta(seconds=30),
    maximum_attempts=3,
)

#5.4 Heartbeat 与长时间 Activity

Python
@activity.defn
async def process_large_dataset(dataset_id: str) -> ProcessResult:
    """处理大数据集,通过 Heartbeat 报告进度"""
    records = await load_dataset(dataset_id)
    processed = 0

    for batch in chunk(records, size=1000):
        result = await process_batch(batch)
        processed += len(batch)

        # 发送 Heartbeat,附带进度信息
        activity.heartbeat({"processed": processed, "total": len(records)})

    return ProcessResult(total_processed=processed)

# Workflow 中配置 Heartbeat 超时
result = await workflow.execute_activity(
    process_large_dataset,
    dataset_id,
    start_to_close_timeout=timedelta(hours=2),
    heartbeat_timeout=timedelta(seconds=30),  # 30 秒无心跳则认为 Activity 失败
    retry_policy=STANDARD_RETRY,
)

#6. Saga 补偿模式

#6.1 什么是 Saga?

在 coomia-dip 的 Action 执行链中,一个 Action 可能涉及多个步骤:

Code
创建订单 → 扣减库存 → 发送通知 → 更新审计日志
    ✅         ✅         ❌ (失败!)
                         ↓
              需要回滚:恢复库存 → 取消订单

#6.2 coomia-dip 的 Saga 实现

Python
from dataclasses import dataclass, field
from typing import Any, Callable, Coroutine

@dataclass
class SagaStep:
    """Saga 的一个步骤"""
    name: str
    action: str  # Activity 函数名
    compensation: str  # 补偿 Activity 函数名
    args: dict = field(default_factory=dict)


@workflow.defn
class SagaWorkflow:
    """通用 Saga 工作流"""

    @workflow.run
    async def run(self, steps: list[SagaStep]) -> dict:
        compensations: list[tuple[str, dict]] = []

        for step in steps:
            try:
                # 执行正向操作
                result = await workflow.execute_activity(
                    step.action,
                    step.args,
                    start_to_close_timeout=timedelta(minutes=5),
                    retry_policy=STANDARD_RETRY,
                )
                # 记录补偿操作(后进先出)
                compensations.append((step.compensation, {
                    **step.args,
                    "forward_result": result,
                }))

            except Exception as e:
                workflow.logger.error(
                    f"Step '{step.name}' failed: {e}. "
                    f"Starting compensation for {len(compensations)} completed steps."
                )
                # 反向执行所有补偿
                await self._compensate(compensations)
                raise

        return {"status": "completed", "steps": len(steps)}

    async def _compensate(self, compensations: list[tuple[str, dict]]) -> None:
        """反向执行补偿操作"""
        errors = []
        for comp_activity, comp_args in reversed(compensations):
            try:
                await workflow.execute_activity(
                    comp_activity,
                    comp_args,
                    start_to_close_timeout=timedelta(minutes=5),
                    retry_policy=RetryPolicy(
                        initial_interval=timedelta(seconds=2),
                        maximum_attempts=5,
                    ),
                )
            except Exception as e:
                # 补偿失败记录但继续执行剩余补偿
                errors.append(f"Compensation '{comp_activity}' failed: {e}")
                workflow.logger.error(f"Compensation failed: {e}")

        if errors:
            raise CompensationError(errors)

#6.3 Action 执行链的 Saga 示例

Python
@workflow.defn
class ActionExecutionWorkflow:
    """coomia-dip Action 执行工作流,带 Saga 补偿"""

    @workflow.run
    async def run(self, action_request: ActionRequest) -> ActionResult:
        compensations = []

        try:
            # Step 1: 验证 Action 参数
            await workflow.execute_activity(
                validate_action_params,
                action_request,
                start_to_close_timeout=timedelta(seconds=30),
            )
            # 验证无需补偿

            # Step 2: 获取乐观锁
            lock = await workflow.execute_activity(
                acquire_lock,
                action_request.target_object_id,
                start_to_close_timeout=timedelta(seconds=10),
            )
            compensations.append(("release_lock", lock))

            # Step 3: 执行数据变更
            mutation_result = await workflow.execute_activity(
                apply_mutation,
                action_request,
                start_to_close_timeout=timedelta(minutes=2),
            )
            compensations.append(("rollback_mutation", mutation_result))

            # Step 4: 触发 Webhook
            webhook_result = await workflow.execute_activity(
                trigger_webhooks,
                action_request,
                start_to_close_timeout=timedelta(seconds=30),
            )
            # Webhook 通知无需补偿(幂等)

            # Step 5: 记录审计日志
            await workflow.execute_activity(
                record_audit,
                action_request,
                start_to_close_timeout=timedelta(seconds=10),
            )

            # Step 6: 释放锁
            await workflow.execute_activity(
                release_lock,
                lock,
                start_to_close_timeout=timedelta(seconds=10),
            )
            compensations = [c for c in compensations if c[0] != "release_lock"]

            return ActionResult(status="success", data=mutation_result)

        except Exception as e:
            # Saga 补偿
            for comp_name, comp_args in reversed(compensations):
                try:
                    await workflow.execute_activity(
                        comp_name,
                        comp_args,
                        start_to_close_timeout=timedelta(minutes=1),
                        retry_policy=RetryPolicy(maximum_attempts=3),
                    )
                except Exception as comp_error:
                    workflow.logger.error(
                        f"Compensation {comp_name} failed: {comp_error}"
                    )

            return ActionResult(status="failed", error=str(e))

#7. Signal 与 Query 通信模式

#7.1 Signal:向运行中的 Workflow 发送异步消息

Python
@workflow.defn
class ApprovalWorkflow:
    def __init__(self):
        self._approval_decision: str | None = None
        self._approval_comment: str = ""

    @workflow.signal
    async def approve(self, comment: str = "") -> None:
        """审批通过信号"""
        self._approval_decision = "approved"
        self._approval_comment = comment

    @workflow.signal
    async def reject(self, comment: str = "") -> None:
        """审批拒绝信号"""
        self._approval_decision = "rejected"
        self._approval_comment = comment

    @workflow.run
    async def run(self, request: ApprovalRequest) -> ApprovalResult:
        # 通知审批人
        await workflow.execute_activity(
            send_approval_notification,
            request,
            start_to_close_timeout=timedelta(seconds=30),
        )

        # 等待审批信号(最多 24 小时)
        try:
            await workflow.wait_condition(
                lambda: self._approval_decision is not None,
                timeout=timedelta(hours=24),
            )
        except asyncio.TimeoutError:
            return ApprovalResult(status="timeout")

        return ApprovalResult(
            status=self._approval_decision,
            comment=self._approval_comment,
        )

#7.2 Query:查询运行中的 Workflow 状态

Python
@workflow.defn
class PipelineWorkflow:
    def __init__(self):
        self._progress = 0
        self._current_stage = "initializing"
        self._error_count = 0

    @workflow.query
    def get_progress(self) -> dict:
        """查询管道执行进度"""
        return {
            "progress": self._progress,
            "stage": self._current_stage,
            "errors": self._error_count,
        }

    @workflow.run
    async def run(self, pipeline_config: PipelineConfig) -> PipelineResult:
        stages = ["extract", "transform", "validate", "load"]

        for i, stage in enumerate(stages):
            self._current_stage = stage
            self._progress = int((i / len(stages)) * 100)

            try:
                await workflow.execute_activity(
                    f"execute_{stage}",
                    pipeline_config,
                    start_to_close_timeout=timedelta(minutes=30),
                )
            except Exception as e:
                self._error_count += 1
                raise

        self._progress = 100
        self._current_stage = "completed"
        return PipelineResult(status="success")

#7.3 从外部发送 Signal / Query

Python
# 从 FastAPI 端点发送 Signal
@app.post("/api/v1/actions/{action_id}/approve")
async def approve_action(action_id: str, body: ApprovalBody):
    client = await Client.connect("temporal-server:7233")

    handle = client.get_workflow_handle(f"action-approval-{action_id}")

    # 发送 Signal
    await handle.signal(ApprovalWorkflow.approve, body.comment)

    return {"status": "signal_sent"}


# 查询 Workflow 状态
@app.get("/api/v1/pipelines/{pipeline_id}/progress")
async def get_pipeline_progress(pipeline_id: str):
    client = await Client.connect("temporal-server:7233")

    handle = client.get_workflow_handle(f"pipeline-{pipeline_id}")

    # Query(同步,立即返回)
    progress = await handle.query(PipelineWorkflow.get_progress)

    return progress

#8. Child Workflow 编排模式

#8.1 扇出模式(Fan-Out)

Python
@workflow.defn
class BatchProcessingWorkflow:
    @workflow.run
    async def run(self, batch_config: BatchConfig) -> BatchResult:
        # 将大批次拆分为多个子批次
        sub_batches = split_into_sub_batches(batch_config, chunk_size=1000)

        # 并行启动 Child Workflow
        handles = []
        for i, sub_batch in enumerate(sub_batches):
            handle = await workflow.start_child_workflow(
                SubBatchWorkflow.run,
                sub_batch,
                id=f"sub-batch-{batch_config.id}-{i}",
                parent_close_policy=ParentClosePolicy.REQUEST_CANCEL,
            )
            handles.append(handle)

        # 等待所有 Child Workflow 完成
        results = await asyncio.gather(
            *[h.result() for h in handles],
            return_exceptions=True,
        )

        # 汇总结果
        successful = sum(1 for r in results if not isinstance(r, Exception))
        failed = sum(1 for r in results if isinstance(r, Exception))

        return BatchResult(
            total=len(sub_batches),
            successful=successful,
            failed=failed,
        )

#8.2 管道模式(Pipeline)

Python
@workflow.defn
class ETLPipelineWorkflow:
    @workflow.run
    async def run(self, config: ETLConfig) -> ETLResult:
        # Stage 1: Extract
        extracted = await workflow.execute_child_workflow(
            ExtractWorkflow.run,
            config.source,
            id=f"etl-extract-{config.id}",
        )

        # Stage 2: Transform(依赖 Extract 结果)
        transformed = await workflow.execute_child_workflow(
            TransformWorkflow.run,
            TransformInput(data=extracted, rules=config.transform_rules),
            id=f"etl-transform-{config.id}",
        )

        # Stage 3: Load(依赖 Transform 结果)
        loaded = await workflow.execute_child_workflow(
            LoadWorkflow.run,
            LoadInput(data=transformed, target=config.target),
            id=f"etl-load-{config.id}",
        )

        return ETLResult(
            records_extracted=extracted.count,
            records_transformed=transformed.count,
            records_loaded=loaded.count,
        )

#9. coomia-dip 6 种工作流模板

模板Task Queue触发方式持续时间补偿
ActionApprovalaction-approvalAPI 调用分钟-天Saga
BatchImportbatch-importSchedule小时Saga
PipelineETLpipeline-etlSchedule/Event分钟-小时Saga
AgentReasoningagent-reasoningEvent秒-分钟
DerivedPropertyCalcderived-propEvent/CDC
DataSyncdata-sync持续运行无限幂等重试

#10. 生产部署与调优

#10.1 Temporal Server 资源规划

组件实例数CPU内存磁盘
Frontend22 核4 GB
History34 核8 GB
Matching22 核4 GB
PostgreSQL1 (HA)8 核32 GBSSD 500 GB
Worker3-54 核8 GB

#10.2 Worker 调优参数

Python
worker = Worker(
    client,
    task_queue="coomia-dip-action-queue",
    workflows=[ActionApprovalWorkflow],
    activities=[validate_action, execute_action],
    # 并发控制
    max_concurrent_workflow_tasks=100,    # 同时处理的 Workflow Task
    max_concurrent_activities=50,         # 同时执行的 Activity
    max_cached_workflows=500,             # 缓存的 Workflow 实例(减少 Replay)
    # Sticky Queue(减少 Replay 频率)
    max_concurrent_workflow_task_polls=5,
    max_concurrent_activity_task_polls=5,
)

#10.3 监控与告警

Python
# 关键指标
TEMPORAL_METRICS = {
    "workflow_task_schedule_to_start_latency":  "Workflow Task 调度延迟",
    "activity_schedule_to_start_latency":       "Activity 调度延迟",
    "workflow_endtoend_latency":                "Workflow 端到端延迟",
    "workflow_failed":                          "Workflow 失败数",
    "activity_execution_failed":                "Activity 执行失败数",
}

#11. Key Takeaways

主题关键结论
架构模型Event Sourcing + Replay 实现持久化执行
DeterminismWorkflow 代码必须确定性,I/O 放 Activity
重试策略指数退避 + 最大次数 + 不可重试错误白名单
Saga 补偿多步操作必须实现补偿,后进先出执行
Signal/QuerySignal 异步通信,Query 同步查询状态
Child WorkflowFan-Out 并行和 Pipeline 串行两种编排
Event History长时间工作流使用 Continue-As-New 防止膨胀
Worker 调优控制并发数 + 启用 Sticky Queue 减少 Replay

下一篇预告:S8-09 将继续深入 Temporal(Part 2),探讨 Schedule、Visibility、Interceptor、多集群复制等高级主题。