Temporal Workflow无法后台运行:当前流程结束前无法触发新流程
问题分析与解决方案
核心问题拆解
- Worker与Workflow触发绑定过紧:你把Worker启动和Workflow执行放在同一个异步上下文里,还
await了execute_workflow,导致程序必须等当前Workflow完全跑完,Worker才会关闭,期间无法处理新的Workflow请求。 - Activity1无限阻塞:Activity1里的循环要执行到count等于10000000000,现实中几乎永远完成不了,直接导致整个Workflow卡壳,无法返回响应。
- Django调用逻辑错误:Django默认是同步框架,直接调用异步的Temporal客户端会有线程模型冲突,而且Worker不应该在Django请求里启动,必须作为独立服务运行。
修复步骤
1. 分离Worker服务与Workflow触发代码
Worker需要是长期运行的独立进程,专门处理Workflow和Activity任务,不能每次触发Workflow都启动一次。
独立的Worker启动代码(比如worker.py):
import asyncio import concurrent from datetime import timedelta from temporalio.client import Client from temporalio.worker import Worker from temporalio import activity, workflow async def get_client(): return await Client.connect("localhost:7233") @activity.defn async def activity1(params) -> dict: print("Inside activity 1") # 替换无限循环为实际业务逻辑 return params @activity.defn async def activity2(params) -> dict: print("Inside activity 2") return params @activity.defn async def activity3(params) -> dict: print("Inside activity 3") return params @activity.defn async def activity4(params) -> dict: print("Inside activity 4") return params @workflow.defn class TestTemporalWorkflow: @workflow.run async def run(self, params) -> dict: response_1 = await workflow.execute_activity( activity1, params, start_to_close_timeout=timedelta(seconds=30), ) response_2 = await workflow.execute_activity( activity2, response_1, start_to_close_timeout=timedelta(seconds=30), ) response_3 = await workflow.execute_activity( activity3, response_2, start_to_close_timeout=timedelta(seconds=30), ) response_4 = await workflow.execute_activity( activity4, response_3, start_to_close_timeout=timedelta(seconds=30), ) return response_4 # 必须返回结果,否则Workflow不会标记完成 async def run_worker(): task_queue = "demo-temporal-workflow-queue" client = await get_client() # 启动Worker并长期运行 worker = Worker( client, task_queue=task_queue, workflows=[TestTemporalWorkflow], activities=[activity1, activity2, activity3, activity4], workflow_task_executor=concurrent.futures.ThreadPoolExecutor(max_workers=10), ) await worker.run() if __name__ == "__main__": asyncio.run(run_worker())
2. 修复Activity1的阻塞问题
删除或替换Activity1里的无限循环,设置合理的业务逻辑终止条件,避免Workflow卡壳。
3. Django中触发Workflow的正确方式
Django默认是同步环境,需用asgiref.sync.async_to_sync将异步的Temporal调用转为同步,同时确保Worker已在后台独立运行。
Django视图示例:
from django.http import JsonResponse from asgiref.sync import async_to_sync import random from datetime import timedelta from temporalio.client import Client from your_module import TestTemporalWorkflow async def trigger_workflow(params): client = await Client.connect("localhost:7233") workflow_id = f"demo-workflow-{random.randint(555, 99999)}" # 用start_workflow替代execute_workflow,避免阻塞到流程结束 handle = await client.start_workflow( TestTemporalWorkflow.run, params, id=workflow_id, task_queue="demo-temporal-workflow-queue", execution_timeout=timedelta(minutes=10), ) return {"workflow_id": workflow_id, "run_id": handle.run_id} def temporal_trigger_view(request): params = {"test": "hello from django"} result = async_to_sync(trigger_workflow)(params) return JsonResponse(result)
4. 关键注意事项
- Worker独立运行:启动Worker进程后保持后台运行,它会持续监听任务队列,处理所有Workflow和Activity请求。
- 优先用
start_workflow:execute_workflow会阻塞到Workflow完成,start_workflow立即返回句柄,更适合后台任务场景。 - Workflow必须返回结果:原代码没有返回值,会导致Temporal认为Workflow仍在运行,必须返回最终结果。
- Django异步兼容:Django 3.1+支持异步视图,可直接编写异步逻辑,无需
async_to_sync。
内容的提问来源于stack exchange,提问作者Vinay Kumar
相关产品推荐
相关产品推荐

