You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Temporal Python SDK中正确实现Worker Interceptor类

修复Temporal Python Workflow拦截器实现问题

正确拦截器实现

针对你的需求(工作流执行前打印名称、保存执行参数),以下是符合Temporal Python SDK规范的拦截器实现,同时解决初始化参数缺失、活动执行失败等问题:

from temporal.workflow import WorkflowInboundInterceptor, WorkflowInterceptorFactory
from temporal.api.common.v1 import WorkflowType

class MyInterceptor(WorkflowInboundInterceptor):
    def __init__(self, next_interceptor):
        super().__init__(next_interceptor)
        self.stored_workflow_params = None  # 用于持久化工作流参数

    async def execute_workflow(self, workflow_type: WorkflowType, input, headers):
        # 需求1:打印工作流名称
        print(f"准备执行工作流: {workflow_type.name}")
        
        # 需求2:保存工作流参数
        self.stored_workflow_params = input
        # 若需全局复用参数,可改为类变量或其他全局存储方式
        # MyInterceptor.global_params = input
        
        # 必须调用下一个拦截器/原生执行逻辑,否则工作流会中断
        return await self.next.execute_workflow(workflow_type, input, headers)

# 拦截器工厂,用于Worker注册
class MyInterceptorFactory(WorkflowInterceptorFactory):
    def create(self, next_interceptor):
        return MyInterceptor(next_interceptor)

完整示例代码

1. 定义活动

from temporal.activity import activity_method

@activity_method(task_queue="my-task-queue", schedule_to_close_timeout_seconds=10)
async def greet_activity(name: str) -> str:
    return f"Hello, {name}!"

2. 定义工作流

from temporal.workflow import workflow_method, Workflow

class GreetWorkflow:
    @workflow_method(task_queue="my-task-queue")
    async def run(self, name: str) -> str:
        return await Workflow.execute_activity(
            greet_activity, 
            name, 
            schedule_to_close_timeout_seconds=10
        )

3. 配置Worker并注册拦截器

from temporal.worker import Worker
from temporal.client import TemporalClient

async def start_worker():
    client = await TemporalClient.connect("localhost:7233")
    worker = Worker(
        client=client,
        task_queue="my-task-queue",
        workflows=[GreetWorkflow],
        activities=[greet_activity],
        # 注册拦截器工厂
        workflow_interceptor_factories=[MyInterceptorFactory()]
    )
    await worker.run()

4. 执行工作流

import asyncio

async def execute_workflow():
    client = await TemporalClient.connect("localhost:7233")
    handle = client.start_workflow(
        GreetWorkflow.run, 
        "Alice", 
        id="greet-workflow-1", 
        task_queue="my-task-queue"
    )
    result = await handle.result()
    print(f"工作流执行结果: {result}")

if __name__ == "__main__":
    # 启动Worker和工作流(实际场景通常分开部署)
    loop = asyncio.get_event_loop()
    tasks = [start_worker(), execute_workflow()]
    loop.run_until_complete(asyncio.gather(*tasks))

错误修复说明

  1. 初始化参数问题:拦截器必须继承WorkflowInboundInterceptor,且__init__方法需接收next_interceptor参数(链式调用要求),之前的错误大概率是缺少该参数。
  2. 活动执行失败:拦截器的execute_workflow方法必须调用self.next.execute_workflow(),否则会中断工作流执行链路,导致活动无法触发。
  3. 参数持久化:通过拦截器实例属性stored_workflow_params保存参数,若需跨实例复用,可改为类变量或外部存储(如内存缓存、数据库)。

内容的提问来源于stack exchange,提问作者Pranbir Sarkar

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.19 03:10:19