如何用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
相关产品推荐
相关产品推荐

