Temporal 工作流引擎深潜(Part 2):Schedule、Visibility、Interceptor 与多集群
1. [Schedule:原生定时调度](#1-schedule原生定时调度)
Coomia发布于 2025年11月17日12 分钟阅读
分享本文Twitter / X
“系列:S8 技术组件深潜 · 第 9 篇 | 难度:高级 | 阅读时间:20 分钟
Temporal 工作流引擎深潜(Part 2):Schedule、Visibility、Interceptor 与多集群
#TL;DR
- 本文是 Temporal 深潜的第二部分,聚焦高级主题:Schedule(替代 Cron 的原生定时调度)、Visibility(基于 Elasticsearch 的工作流检索)、Interceptor(横切关注点注入)、以及多集群复制方案
- 深入分析 Temporal 的 Namespace 隔离模型、Search Attributes 自定义、Worker Versioning 策略、以及 coomia-dip 中的可观测性集成方案
- 包含 coomia-dip 生产环境中的故障排查手册和性能调优指南
#目录
- Schedule:原生定时调度
- Visibility:工作流检索与监控
- Search Attributes 自定义
- Interceptor 拦截器
- Worker Versioning
- Namespace 隔离模型
- 多集群复制
- 可观测性集成
- 故障排查手册
- 性能基准与极限测试
- Key Takeaways
#1. Schedule:原生定时调度
#1.1 为什么不用 Cron?
传统 Cron 方案的痛点:
| 问题 | Cron | Temporal Schedule |
|---|---|---|
| 错过执行(节点宕机) | 丢失 | 自动 Backfill |
| 重叠执行 | 需要手动锁 | 原生 Overlap Policy |
| 执行历史 | 仅日志 | 完整 Event History |
| 暂停/恢复 | 改 crontab | API 一键操作 |
| 参数化调度 | 困难 | 原生支持 |
#1.2 coomia-dip Schedule 配置
Python
from temporalio.client import (
Client,
Schedule,
ScheduleActionStartWorkflow,
ScheduleIntervalSpec,
ScheduleSpec,
ScheduleOverlapPolicy,
ScheduleState,
)
async def create_report_schedule(client: Client, world_id: str) -> str:
"""创建每日报告生成的定时调度"""
schedule_handle = await client.create_schedule(
id=f"daily-report-{world_id}",
schedule=Schedule(
action=ScheduleActionStartWorkflow(
workflow="ReportGenerationWorkflow",
arg={"world_id": world_id, "report_type": "daily"},
id=f"report-{world_id}-{{{{.ScheduledTime.Format `20060102`}}}}",
task_queue="coomia-dip-report-queue",
),
spec=ScheduleSpec(
intervals=[
ScheduleIntervalSpec(
every=timedelta(hours=24),
offset=timedelta(hours=6), # 每天 06:00 UTC 执行
),
],
),
policy=ScheduleOverlapPolicy.SKIP, # 上次未完成则跳过
state=ScheduleState(
note="Daily report for world " + world_id,
limited_actions=False,
),
),
)
return schedule_handle.id
#1.3 Overlap Policy 策略
Python
# 4 种重叠策略
OVERLAP_POLICIES = {
"SKIP": "上次未完成,跳过本次",
"BUFFER_ONE": "缓冲一次,上次完成后立即执行",
"BUFFER_ALL": "缓冲所有,队列执行",
"CANCEL_OTHER": "取消上次,启动新的",
"TERMINATE_OTHER": "终止上次,启动新的",
"ALLOW_ALL": "允许并行执行",
}
# coomia-dip 场景推荐
SCHEDULE_CONFIGS = {
"daily_report": ScheduleOverlapPolicy.SKIP, # 报告可以跳过
"data_sync": ScheduleOverlapPolicy.BUFFER_ONE, # 同步不能跳过
"cleanup": ScheduleOverlapPolicy.SKIP, # 清理可以跳过
"metric_aggregate": ScheduleOverlapPolicy.CANCEL_OTHER, # 用最新数据
}
#1.4 Backfill:补执行错过的调度
Python
async def backfill_missed_schedules(
client: Client,
schedule_id: str,
start: datetime,
end: datetime,
) -> None:
"""补执行因维护窗口错过的调度"""
handle = client.get_schedule_handle(schedule_id)
await handle.backfill(
ScheduleBackfill(
start_at=start,
end_at=end,
overlap=ScheduleOverlapPolicy.BUFFER_ALL,
),
)
#2. Visibility:工作流检索与监控
#2.1 Visibility Store 架构
Code
Temporal Server
│
├── Standard Visibility (PostgreSQL)
│ └── 基础过滤:WorkflowType, Status, StartTime
│
└── Advanced Visibility (Elasticsearch)
└── 全功能:自定义 Search Attributes, 全文搜索, 复杂查询
coomia-dip 使用 Elasticsearch 作为 Visibility Store:
YAML
# Temporal Server 配置
persistence:
advancedVisibilityStore: es-visibility
datastores:
es-visibility:
elasticsearch:
version: v7
url:
scheme: https
host: elasticsearch:9200
indices:
visibility: temporal_visibility_v1
#2.2 List Filter 查询语法
Python
# 查询所有失败的 Action Workflow
workflows = await client.list_workflows(
query='WorkflowType = "ActionExecutionWorkflow" AND ExecutionStatus = "Failed"'
)
# 查询特定 World 的运行中工作流
workflows = await client.list_workflows(
query=(
'CustomStringField = "world-123" '
'AND ExecutionStatus = "Running" '
'AND StartTime > "2026-03-01T00:00:00Z"'
)
)
# 查询超时的审批工作流
workflows = await client.list_workflows(
query=(
'WorkflowType = "ApprovalWorkflow" '
'AND ExecutionStatus = "Running" '
'AND StartTime < "2026-03-20T00:00:00Z"'
)
)
#3. Search Attributes 自定义
#3.1 注册自定义 Search Attributes
Python
# 为 coomia-dip 注册自定义搜索属性
async def register_search_attributes(client: Client) -> None:
await client.operator_service.add_search_attributes(
namespace="coomia-dip-default",
search_attributes={
"WorldId": SearchAttributeType.KEYWORD,
"ObjectType": SearchAttributeType.KEYWORD,
"ActionType": SearchAttributeType.KEYWORD,
"Priority": SearchAttributeType.INT,
"Initiator": SearchAttributeType.KEYWORD,
"ErrorMessage": SearchAttributeType.TEXT,
"DataSize": SearchAttributeType.DOUBLE,
"Tags": SearchAttributeType.KEYWORD_LIST,
},
)
#3.2 在 Workflow 中设置 Search Attributes
Python
@workflow.defn
class ActionExecutionWorkflow:
@workflow.run
async def run(self, request: ActionRequest) -> ActionResult:
# 启动时设置 Search Attributes(通过 start_workflow 参数)
# 运行中动态更新
workflow.upsert_search_attributes(
[
SearchAttributeUpdate(
SearchAttributeKey.for_keyword("WorldId"),
request.world_id,
),
SearchAttributeUpdate(
SearchAttributeKey.for_keyword("ActionType"),
request.action_type,
),
SearchAttributeUpdate(
SearchAttributeKey.for_int("Priority"),
request.priority,
),
]
)
try:
result = await self._execute(request)
return result
except Exception as e:
# 失败时更新错误信息到 Search Attributes
workflow.upsert_search_attributes(
[
SearchAttributeUpdate(
SearchAttributeKey.for_text("ErrorMessage"),
str(e),
),
]
)
raise
#4. Interceptor 拦截器
#4.1 Interceptor 架构
Code
Client Call → Client Interceptor → Temporal Server
↓
Worker Poll ← Activity Interceptor ← Workflow Interceptor
#4.2 日志与追踪 Interceptor
Python
from temporalio.worker import (
Interceptor,
ExecuteWorkflowInput,
ExecuteActivityInput,
)
class TracingInterceptor(Interceptor):
"""coomia-dip 分布式追踪拦截器"""
def intercept_activity(self, next_interceptor):
return TracingActivityInterceptor(next_interceptor)
def workflow_interceptor_class(self, input):
return TracingWorkflowInterceptor
class TracingActivityInterceptor:
def __init__(self, next_interceptor):
self.next = next_interceptor
async def execute_activity(self, input: ExecuteActivityInput):
activity_name = input.fn.__name__
start_time = time.monotonic()
# 注入 trace context
span = tracer.start_span(
f"temporal.activity.{activity_name}",
attributes={
"temporal.workflow_id": activity.info().workflow_id,
"temporal.activity_id": activity.info().activity_id,
"temporal.task_queue": activity.info().task_queue,
"temporal.attempt": activity.info().attempt,
},
)
try:
result = await self.next.execute_activity(input)
span.set_status(StatusCode.OK)
return result
except Exception as e:
span.set_status(StatusCode.ERROR, str(e))
span.record_exception(e)
raise
finally:
duration = time.monotonic() - start_time
metrics.histogram(
"temporal_activity_duration_seconds",
duration,
tags={"activity": activity_name},
)
span.end()
#4.3 审计 Interceptor
Python
class AuditInterceptor(Interceptor):
"""记录所有工作流启动和完成的审计日志"""
def intercept_activity(self, next_interceptor):
return AuditActivityInterceptor(next_interceptor)
def workflow_interceptor_class(self, input):
return AuditWorkflowInterceptor
class AuditWorkflowInterceptor:
async def execute_workflow(self, input: ExecuteWorkflowInput):
workflow_type = type(input.workflow).__name__
workflow_id = workflow.info().workflow_id
# 记录工作流开始
await self._log_audit_event(
event_type="WORKFLOW_STARTED",
workflow_type=workflow_type,
workflow_id=workflow_id,
args=str(input.args)[:500], # 截断避免日志过大
)
try:
result = await input.execute()
await self._log_audit_event(
event_type="WORKFLOW_COMPLETED",
workflow_type=workflow_type,
workflow_id=workflow_id,
)
return result
except Exception as e:
await self._log_audit_event(
event_type="WORKFLOW_FAILED",
workflow_type=workflow_type,
workflow_id=workflow_id,
error=str(e),
)
raise
#5. Worker Versioning
#5.1 Build ID 版本化(Flink 1.24+)
Python
# 注册新 Build ID
async def register_worker_version(
client: Client,
task_queue: str,
build_id: str,
existing_compatible_id: str | None = None,
) -> None:
if existing_compatible_id:
# 新版本与旧版本兼容(相同 Task Queue)
await client.update_worker_build_id_compatibility(
task_queue,
BuildIdOpAddNewCompatible(
new_build_id=build_id,
existing_compatible_build_id=existing_compatible_id,
),
)
else:
# 新版本不兼容旧版本(新 Task Queue Set)
await client.update_worker_build_id_compatibility(
task_queue,
BuildIdOpAddNewDefault(build_id),
)
# Worker 声明自己的 Build ID
worker = Worker(
client,
task_queue="coomia-dip-action-queue",
workflows=[ActionApprovalWorkflow],
activities=[validate_action],
build_id="v2.3.0-abc123",
use_worker_versioning=True,
)
#5.2 版本化部署策略
Code
v1.0 Worker ────────────────────────────────────┐
(处理 v1.0 的运行中工作流) │
│
v2.0 Worker ─────────────────────────────────────┤
(处理新启动的工作流 + v2.0 兼容的旧工作流) │
│
时间 ──────────────────────────────────────────────→
部署 v2.0 v1.0 工作流 关闭 v1.0
全部完成 Worker
#6. Namespace 隔离模型
#6.1 coomia-dip 的 Namespace 策略
Code
coomia-dip Temporal Namespaces:
│
├── coomia-dip-default # 默认 Namespace(开发/测试)
├── coomia-dip-world-{id} # 每个 World 一个 Namespace
│ ├── Task Queue: action-approval
│ ├── Task Queue: pipeline-etl
│ └── Task Queue: agent-reasoning
├── coomia-dip-system # 系统内部工作流
│ ├── Task Queue: maintenance
│ └── Task Queue: monitoring
└── coomia-dip-staging # 预发布环境
#6.2 Namespace 级别的资源限制
Python
# 创建带资源限制的 Namespace
async def create_world_namespace(
client: Client,
world_id: str,
tier: str = "standard",
) -> None:
limits = {
"standard": {
"max_workflow_execution_count": 10_000,
"max_concurrent_workflow_tasks": 200,
"workflow_execution_rate_limit": 100, # 每秒
},
"premium": {
"max_workflow_execution_count": 100_000,
"max_concurrent_workflow_tasks": 1000,
"workflow_execution_rate_limit": 500,
},
}
tier_limits = limits[tier]
await client.operator_service.create_namespace(
name=f"coomia-dip-world-{world_id}",
retention_period=timedelta(days=30),
# 资源限制通过 Temporal Server 配置
)
#7. 多集群复制
#7.1 Multi-Cluster Replication 架构
Code
Region A (Primary) Region B (Standby)
┌──────────────────┐ ┌──────────────────┐
│ Temporal Server │ │ Temporal Server │
│ ┌────────────┐ │ Replication │ ┌────────────┐ │
│ │ History │──┼──────────────┼──│ History │ │
│ │ Service │ │ │ │ Service │ │
│ └────────────┘ │ │ └────────────┘ │
│ ┌────────────┐ │ │ ┌────────────┐ │
│ │ PostgreSQL │──┼──────────────┼──│ PostgreSQL │ │
│ └────────────┘ │ │ └────────────┘ │
└──────────────────┘ └──────────────────┘
↑ ↑
Workers A Workers B
(Active) (Standby)
#7.2 故障切换流程
Python
async def failover_to_standby(client: Client, namespace: str) -> None:
"""将 Namespace 切换到备用集群"""
# 1. 更新 Namespace 的活跃集群
await client.operator_service.update_namespace(
namespace=namespace,
update_info=NamespaceUpdateInfo(
active_cluster_name="region-b",
),
)
# 2. 等待复制完成
await asyncio.sleep(5)
# 3. 验证切换成功
info = await client.operator_service.describe_namespace(namespace)
assert info.active_cluster_name == "region-b"
#8. 可观测性集成
#8.1 Metrics 集成
Python
from temporalio.runtime import PrometheusConfig, Runtime, TelemetryConfig
# 配置 Prometheus Metrics
runtime = Runtime(
telemetry=TelemetryConfig(
metrics=PrometheusConfig(
bind_address="0.0.0.0:9464",
),
),
)
client = await Client.connect(
"temporal-server:7233",
runtime=runtime,
)
#8.2 关键 Grafana Dashboard
YAML
# coomia-dip Temporal Grafana Dashboard 核心面板
panels:
- title: "Workflow 启动速率"
query: rate(temporal_workflow_started_total[5m])
- title: "Workflow 失败率"
query: |
rate(temporal_workflow_failed_total[5m]) /
rate(temporal_workflow_completed_total[5m])
- title: "Activity 延迟分布"
query: histogram_quantile(0.99, temporal_activity_execution_latency_bucket)
- title: "Task Queue 积压"
query: temporal_task_queue_backlog
- title: "Schedule 错过次数"
query: increase(temporal_schedule_missed_catchup_window_total[1h])
- title: "Worker Slot 使用率"
query: |
temporal_worker_task_slots_used /
temporal_worker_task_slots_available
#8.3 OpenTelemetry 集成
Python
from opentelemetry import trace
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
from temporalio.contrib.opentelemetry import TracingInterceptor
# 配置 OpenTelemetry
tracer_provider = TracerProvider(
resource=Resource.create({"service.name": "coomia-dip-worker"}),
)
tracer_provider.add_span_processor(
BatchSpanProcessor(OTLPSpanExporter(endpoint="otel-collector:4317"))
)
trace.set_tracer_provider(tracer_provider)
# 使用 Temporal 的 OpenTelemetry Interceptor
client = await Client.connect(
"temporal-server:7233",
interceptors=[TracingInterceptor()],
)
worker = Worker(
client,
task_queue="coomia-dip-action-queue",
workflows=[ActionApprovalWorkflow],
activities=[validate_action],
interceptors=[TracingInterceptor()],
)
#9. 故障排查手册
#9.1 常见问题与解决方案
| 问题 | 症状 | 排查步骤 | 解决方案 |
|---|---|---|---|
| Workflow Task 卡住 | 工作流长时间 Running | 检查 Worker 连接、Task Queue 名称 | 重启 Worker 或修复 Task Queue |
| Activity 超时 | 频繁重试 | 检查 Heartbeat、外部依赖 | 增大超时时间或优化 Activity |
| Non-Determinism Error | Workflow Task 失败 | 查看 Event History 对比 | 使用 Patching API 修复 |
| Schedule 不执行 | 定时任务未触发 | 检查 Schedule 状态(Paused?) | 恢复 Schedule |
| 状态大小超限 | ContinueAsNew 后失败 | 检查传递的 State 大小 | 减少传递数据,存外部存储 |
#9.2 调试命令
Bash
# 查看 Workflow 执行详情
temporal workflow describe --workflow-id action-approval-123
# 查看 Event History
temporal workflow show --workflow-id action-approval-123
# 列出运行中工作流
temporal workflow list --query 'ExecutionStatus="Running"'
# 发送 Signal
temporal workflow signal --workflow-id action-approval-123 \
--name approve --input '"approved by admin"'
# 取消工作流
temporal workflow cancel --workflow-id action-approval-123
# 终止工作流(紧急)
temporal workflow terminate --workflow-id action-approval-123 \
--reason "Manual termination for debugging"
# 查看 Task Queue 状态
temporal task-queue describe --task-queue coomia-dip-action-queue
# 查看 Schedule 状态
temporal schedule describe --schedule-id daily-report-world-123
#9.3 Non-Determinism 调试
Python
# 启用 Replay 测试(单元测试中验证 Determinism)
from temporalio.testing import WorkflowEnvironment
from temporalio.worker import Replayer
async def test_workflow_replay():
"""验证 Workflow 代码修改不破坏 Determinism"""
# 从生产环境导出 Event History
history_json = await export_workflow_history(
workflow_id="action-approval-prod-123"
)
# Replay 测试
replayer = Replayer(workflows=[ActionApprovalWorkflow])
try:
await replayer.replay_workflow(
WorkflowHistory.from_json("action-approval-prod-123", history_json)
)
print("Replay succeeded: Workflow is deterministic")
except Exception as e:
print(f"Replay failed: Non-determinism detected: {e}")
#10. 性能基准与极限测试
#10.1 单集群吞吐量
| 测试场景 | 工作流/秒 | Activity/秒 | P99 延迟 |
|---|---|---|---|
| 简单工作流(1 Activity) | 2,000 | 2,000 | 50 ms |
| 中等工作流(5 Activity) | 800 | 4,000 | 200 ms |
| 复杂工作流(10 Activity + Timer) | 300 | 3,000 | 500 ms |
| Saga 工作流(5 步 + 补偿) | 500 | 5,000 | 300 ms |
#10.2 资源瓶颈分析
Code
工作流吞吐量受限因素:
1. History Service CPU → 增加实例数
2. PostgreSQL IOPS → 升级存储 / 读写分离
3. Elasticsearch 写入延迟 → 增加 Shard / 升级节点
4. Worker 并发能力 → 水平扩展 Worker
5. 网络带宽 → 升级网络
#10.3 coomia-dip 容量规划公式
Code
所需 History 实例 = ceil(目标工作流QPS × 平均事件数 / 单实例处理能力)
示例:
目标:500 QPS,平均 20 个事件/工作流
单实例处理能力:3000 事件/秒
所需实例 = ceil(500 × 20 / 3000) = ceil(3.33) = 4 实例
所需 Worker 实例 = ceil(目标 Activity QPS / 单 Worker Activity 并发)
示例:
目标:2000 Activity/秒,单 Worker 50 并发
所需 Worker = ceil(2000 / 50) = 40 Worker
考虑冗余 × 1.5 = 60 Worker
#11. Key Takeaways
| 主题 | 关键结论 |
|---|---|
| Schedule | 替代 Cron,支持 Backfill 和 Overlap Policy |
| Visibility | Elasticsearch 驱动,自定义 Search Attributes |
| Interceptor | 用于追踪、审计、限流等横切关注点 |
| Worker Versioning | Build ID 实现安全的滚动升级 |
| Namespace | 每个 World 一个 Namespace 实现隔离 |
| 多集群 | 跨 Region 复制支持灾备切换 |
| 可观测性 | Prometheus + Grafana + OpenTelemetry 三位一体 |
| 容量规划 | 基于 QPS 和事件数推导实例数 |
“下一篇预告:S8-10 将深入 DolphinScheduler,探讨 coomia-dip 如何用 DolphinScheduler 实现批处理 DAG 调度与数据管道编排。