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

如何用asyncio执行DAG任务?解决现有实现的并行瓶颈

问题分析与解决方案

问题原因

你的代码中,execute_dag 每次会一次性取出 ready_queue 里的所有任务,存入 current_tasks 后用 await asyncio.gather(*current_tasks) 等待这批任务全部完成,才会进入下一轮循环处理新的就绪任务。

在你的DAG结构里,任务2和任务5属于同一批启动的任务。当任务2完成时,任务3的前置条件(仅依赖任务2)已经满足,会被加入 ready_queue,但此时当前批次的 gather 仍在等待任务5完成,下一轮循环要等任务5结束后才会执行,因此任务3被迫等待任务5完成才能启动,造成了不必要的并行瓶颈。

解决办法

要实现任务的最大化并行,核心是避免批量等待全部任务完成,改为在单个任务完成后立即检查并启动可执行的后继任务,同时持续处理就绪队列。以下是两种可行的修改方案:

方案一:直接启动后继任务(移除就绪队列)

async def execute_task(self, task_id):
    await self.tasks[task_id]()
    self.task_status[task_id] = "done"

    async with self.lock:
        # 任务完成后立即检查并启动所有满足条件的后继任务
        for successor in self.dag.successors(task_id):
            if all(self.task_status[predecessor] == "done" for predecessor in self.dag.predecessors(successor)):
                self.task_status[successor] = "running"
                # 直接创建任务并加入活跃任务集合
                self.active_tasks.add(asyncio.create_task(self.execute_task(successor)))

async def execute_dag(self):
    self.active_tasks = set()
    # 初始化启动所有入度为0的根任务
    async with self.lock:
        for task in self.dag.nodes:
            if self.dag.in_degree(task) == 0:
                self.task_status[task] = "running"
                self.active_tasks.add(asyncio.create_task(self.execute_task(task)))
    
    # 持续等待任意任务完成,直到所有任务执行完毕
    while self.active_tasks:
        done, pending = await asyncio.wait(self.active_tasks, return_when=asyncio.FIRST_COMPLETED)
        self.active_tasks.difference_update(done)

方案二:保留就绪队列,持续检查并启动

async def execute_task(self, task_id):
    await self.tasks[task_id]()
    self.task_status[task_id] = "done"

    async with self.lock:
        for successor in self.dag.successors(task_id):
            if all(self.task_status[predecessor] == "done" for predecessor in self.dag.predecessors(successor)):
                self.ready_queue.append(successor)

async def execute_dag(self):
    self.active_tasks = set()
    # 初始化就绪队列
    for task in self.dag.nodes:
        if self.dag.in_degree(task) == 0:
            self.ready_queue.append(task)
    
    # 持续处理就绪队列和活跃任务
    while self.ready_queue or self.active_tasks:
        # 启动当前所有就绪任务
        async with self.lock:
            while self.ready_queue:
                task_id = self.ready_queue.popleft()
                self.task_status[task_id] = "running"
                self.active_tasks.add(asyncio.create_task(self.execute_task(task_id)))
        
        # 等待任意一个任务完成,再循环检查新的就绪任务
        if self.active_tasks:
            done, pending = await asyncio.wait(self.active_tasks, return_when=asyncio.FIRST_COMPLETED)
            self.active_tasks.difference_update(done)

两种方案的核心逻辑一致:不再等待一批任务全部完成,而是在单个任务结束后立即触发后继任务的检查与启动,确保满足条件的任务能尽快并行执行,避免不必要的等待。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 00:45:17