返回博客

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 生产环境中的故障排查手册和性能调优指南

#目录

  1. Schedule:原生定时调度
  2. Visibility:工作流检索与监控
  3. Search Attributes 自定义
  4. Interceptor 拦截器
  5. Worker Versioning
  6. Namespace 隔离模型
  7. 多集群复制
  8. 可观测性集成
  9. 故障排查手册
  10. 性能基准与极限测试
  11. Key Takeaways

#1. Schedule:原生定时调度

#1.1 为什么不用 Cron?

传统 Cron 方案的痛点:

问题CronTemporal Schedule
错过执行(节点宕机)丢失自动 Backfill
重叠执行需要手动锁原生 Overlap Policy
执行历史仅日志完整 Event History
暂停/恢复改 crontabAPI 一键操作
参数化调度困难原生支持

#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

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 ErrorWorkflow 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,0002,00050 ms
中等工作流(5 Activity)8004,000200 ms
复杂工作流(10 Activity + Timer)3003,000500 ms
Saga 工作流(5 步 + 补偿)5005,000300 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
VisibilityElasticsearch 驱动,自定义 Search Attributes
Interceptor用于追踪、审计、限流等横切关注点
Worker VersioningBuild ID 实现安全的滚动升级
Namespace每个 World 一个 Namespace 实现隔离
多集群跨 Region 复制支持灾备切换
可观测性Prometheus + Grafana + OpenTelemetry 三位一体
容量规划基于 QPS 和事件数推导实例数

下一篇预告:S8-10 将深入 DolphinScheduler,探讨 coomia-dip 如何用 DolphinScheduler 实现批处理 DAG 调度与数据管道编排。