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 种典型工作流模板
#目录
- Temporal 在 coomia-dip 中的定位
- Temporal Server 四组件架构
- Event Sourcing 持久化模型
- Workflow Determinism 约束
- Activity 重试策略的数学模型
- Saga 补偿模式
- Signal 与 Query 通信模式
- Child Workflow 编排模式
- coomia-dip 6 种工作流模板
- 生产部署与调优
- 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 的关键规则:
- Workflow 代码的决策路径在 Replay 时必须产生相同结果(Determinism)
- Activity 的实际执行不会重复 — 使用 Event History 中记录的结果
- 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,000 | 50,000 |
| 最大 History 大小 | 50 MB | 50 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 | 触发方式 | 持续时间 | 补偿 |
|---|---|---|---|---|
| ActionApproval | action-approval | API 调用 | 分钟-天 | Saga |
| BatchImport | batch-import | Schedule | 小时 | Saga |
| PipelineETL | pipeline-etl | Schedule/Event | 分钟-小时 | Saga |
| AgentReasoning | agent-reasoning | Event | 秒-分钟 | 无 |
| DerivedPropertyCalc | derived-prop | Event/CDC | 秒 | 无 |
| DataSync | data-sync | 持续运行 | 无限 | 幂等重试 |
#10. 生产部署与调优
#10.1 Temporal Server 资源规划
| 组件 | 实例数 | CPU | 内存 | 磁盘 |
|---|---|---|---|---|
| Frontend | 2 | 2 核 | 4 GB | 无 |
| History | 3 | 4 核 | 8 GB | 无 |
| Matching | 2 | 2 核 | 4 GB | 无 |
| PostgreSQL | 1 (HA) | 8 核 | 32 GB | SSD 500 GB |
| Worker | 3-5 | 4 核 | 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 实现持久化执行 |
| Determinism | Workflow 代码必须确定性,I/O 放 Activity |
| 重试策略 | 指数退避 + 最大次数 + 不可重试错误白名单 |
| Saga 补偿 | 多步操作必须实现补偿,后进先出执行 |
| Signal/Query | Signal 异步通信,Query 同步查询状态 |
| Child Workflow | Fan-Out 并行和 Pipeline 串行两种编排 |
| Event History | 长时间工作流使用 Continue-As-New 防止膨胀 |
| Worker 调优 | 控制并发数 + 启用 Sticky Queue 减少 Replay |
“下一篇预告:S8-09 将继续深入 Temporal(Part 2),探讨 Schedule、Visibility、Interceptor、多集群复制等高级主题。