单Worker实例中Workflow与Activity的并发执行问题
单Temporal Worker实例中多个Workflow串行执行,如何实现并行?
我们在集群中以单副本Pod运行Temporal Worker(不希望多副本),期望同一Worker实例中并行执行多个Workflow,无需等待当前Workflow完成即可启动新的Workflow。已在Worker端配置了Workflow和任务的最大并发参数,但目前多个Workflow仍串行执行,总耗时是预期的两倍。
代码与配置
def slp(sec): sleep(sec) return f"slept {sec} sec" @activity.defn(name="Sleeping 1") async def sleeping1(): response = slp(20) return response @activity.defn(name="sleeping 2") async def sleeping2(): response = slp(120) return response @workflow.defn class sleepingWF: @workflow.run async def run(self, body): s1 = await workflow.execute_activity(sleeping1) s2 = await workflow.execute_activity(sleeping2) async def main(): TEMPORAL_ENDPOINT = "localhost:7233" TASK_QUEUE = "mymac" client = await Client.connect(target_host=TEMPORAL_ENDPOINT) with concurrent.futures.ThreadPoolExecutor(max_workers=100) as activity_executor: worker = Worker( client, task_queue=TASK_QUEUE, workflows=[sleepingWF], activities=[sleeping1, sleeping2], activity_executor=activity_executor, max_concurrent_workflow_tasks=100, max_concurrent_activities=100, max_concurrent_workflow_task_polls=10, max_concurrent_activity_task_polls=10, ) await worker.run() if __name__ == "__main__": asyncio.run(main())
期望结果
workflow_execution_1 : start_time -> 00:00:00 end_time -> 00:02:20 workflow_execution_2 : start_time -> 00:00:00 end_time -> 00:02:20
两次执行总耗时应为140秒。
当前实际情况
workflow_execution_1 : start_time -> 00:00:00 end_time -> 00:02:20 workflow_execution_2 : start_time -> 00:02:20 end_time -> 00:04:40
两次执行总耗时为280秒。
环境
Python 3.11
问题根源
代码中存在两个关键问题导致Workflow串行执行:
- 同步阻塞函数阻塞事件循环:
slp函数使用了同步sleep(sec),在异步Activity函数中直接调用同步阻塞操作会卡住整个asyncio事件循环,导致Worker无法在阻塞期间处理其他Workflow或Activity任务。 - Activity未正确利用线程池:虽然配置了
activity_executor为ThreadPoolExecutor,但异步Activity函数并未将同步任务提交到该线程池执行,而是直接在事件循环线程中运行同步代码。
修复方案
方案1:将同步sleep改为异步sleep(推荐)
如果业务逻辑可改为异步实现,直接用asyncio.sleep替代同步sleep,避免阻塞事件循环:
import asyncio async def slp(sec): await asyncio.sleep(sec) return f"slept {sec} sec" @activity.defn(name="Sleeping 1") async def sleeping1(): response = await slp(20) return response @activity.defn(name="sleeping 2") async def sleeping2(): response = await slp(120) return response
方案2:将同步任务提交到线程池执行
如果必须使用同步阻塞函数,在异步Activity中通过loop.run_in_executor将同步任务提交到配置的线程池:
import asyncio def slp(sec): sleep(sec) return f"slept {sec} sec" @activity.defn(name="Sleeping 1") async def sleeping1(): loop = asyncio.get_running_loop() response = await loop.run_in_executor(activity_executor, slp, 20) return response @activity.defn(name="sleeping 2") async def sleeping2(): loop = asyncio.get_running_loop() response = await loop.run_in_executor(activity_executor, slp, 120) return response
额外优化:Workflow内部Activity并行执行
若Workflow内部的两个Activity无需串行,可改为并行调用缩短单个Workflow的执行时间:
@workflow.defn class sleepingWF: @workflow.run async def run(self, body): # 并行执行两个Activity s1_future = workflow.execute_activity(sleeping1) s2_future = workflow.execute_activity(sleeping2) s1, s2 = await asyncio.gather(s1_future, s2_future)
验证效果
修复后,两个Workflow实例可同时启动,各自的Activity并行执行,总耗时将符合预期的140秒左右。
内容的提问来源于stack exchange,提问作者sriram manikanth
相关产品推荐
相关产品推荐

