如何在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))
错误修复说明
- 初始化参数问题:拦截器必须继承
WorkflowInboundInterceptor,且__init__方法需接收next_interceptor参数(链式调用要求),之前的错误大概率是缺少该参数。 - 活动执行失败:拦截器的
execute_workflow方法必须调用self.next.execute_workflow(),否则会中断工作流执行链路,导致活动无法触发。 - 参数持久化:通过拦截器实例属性
stored_workflow_params保存参数,若需跨实例复用,可改为类变量或外部存储(如内存缓存、数据库)。
内容的提问来源于stack exchange,提问作者Pranbir Sarkar
相关产品推荐
相关产品推荐

