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

单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串行执行:

  1. 同步阻塞函数阻塞事件循环:slp函数使用了同步sleep(sec),在异步Activity函数中直接调用同步阻塞操作会卡住整个asyncio事件循环,导致Worker无法在阻塞期间处理其他Workflow或Activity任务。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 13:35:56