asyncio.gather无法处理任务依赖,需用task+依赖字典实现dag调度:先create_task预注册,再按拓扑序await,避免重复await同一task,并确保异常传播与取消传递。

asyncio.gather 不能直接处理依赖关系
很多人一上来就用 asyncio.gather 并发跑一堆协程,结果发现任务 A 依赖 B 的返回值,但 B 还没完成 A 就开始执行了——gather 只保证并发启动,不保证执行顺序或数据流。它本质是“同时发车”,不是“流水线”。
真正需要的是:B 完成后自动触发 A,A 完成后再触发 C,且整个过程保持异步、不阻塞事件循环。
- 错误做法:在 A 里
await B()—— 这会让 A 等待 B,但 B 本身可能还没被调度,变成串行,失去并发价值 - 正确思路:把依赖关系建模为有向图,按拓扑序调度,但每个节点仍是 async 函数,靠
asyncio.create_task提前注册,再用await显式等待其完成 - 关键约束:不能让下游任务重复 await 同一个 task(否则会报
RuntimeError: cannot reuse already awaited coroutine)
用 asyncio.Task + 依赖字典实现轻量级 DAG 调度
不需要引入 aiodag 或 prefect 这类重型框架,几行代码就能搭出可读、可调试的依赖调度器。核心是维护一个 task_map: Dict[str, asyncio.Task],把任务名映射到已创建的 Task 实例。
示例场景:任务 "download" 输出路径,"parse" 需要该路径,"report" 需要 parse 的结构化结果:
deps = {
"parse": ["download"],
"report": ["parse"]
}
调度逻辑要点:
SkillSub Pro - Python 题解与代码注释双功能技能功能概述SkillSub Pro - Python 题解与代码注释双功能技能是一项面向实际任务的技能,主要用于SkillSub Pro 是一个 Python 题解生成与代码注释的 双功能合体技能 ,专为学生、算法学习者和开发者设计;✅ 一个技能,两种用途 :;核心要点📝 题解模式 :输入题目/题号,自动生成完整 Python 题解(含详细注释、解题思路、复杂度分析);💬 注释模式 :输入 Python 代码,自动添加详细中。它将相关步骤、
- 先遍历所有任务,对每个任务调用
asyncio.create_task(func()),但**暂不 await**,存入task_map - 再按拓扑序(如用
networkx.algorithms.dag.topological_sort,或手写 Kahn 算法)依次处理任务名 - 对每个任务名
name,先await task_map[name],确保它及其所有前置任务都已完成 - 注意:如果某任务被多个下游依赖,只需 await 一次;重复 await 同一个
Task对象会抛错
避免 await 已完成 task 导致的 RuntimeError
常见错误是把 task_map[name] 当作协程反复传给 await,比如写成 await parse_task()(函数调用)或误把 task_map[name] 当协程对象重新 await。Task 对象只能 await 一次。
安全做法:
- 始终用
task_map[name]存储asyncio.Task实例,不是协程对象 - 检查是否已完成:用
task_map[name].done()判断,但不要据此决定是否 await —— 即使 done,await 仍合法(返回缓存结果),但若已 cancel 或异常,await 会重新抛出 - 更稳妥:统一用
await task_map[name],并在外层 try/except 捕获可能的异常,而不是跳过 await - 切勿写
await task_map[name]();——Task不可调用,这是典型类型混淆
复杂依赖中容易忽略的异常传播与取消传递
当 "download" 失败时,"parse" 不应静默跳过,而应收到原始异常;同样,手动 cancel() 顶层调度器,所有下游 task 应响应取消信号。
这要求:
- 每个 task 必须在异常时显式 re-raise,不能吞掉 ——
asyncio.Task.exception()返回None表示未完成或无异常,但 await 才真正触发传播 - 调度循环中,一旦某个 task await 抛出异常,后续依赖 task 的 await 应立即中断(可用
asyncio.shield包裹,但更推荐让异常自然冒泡) - 若需支持外部取消,应在调度主协程中监听
asyncio.current_task().cancelled(),并主动调用task.cancel()清理未完成 task - 特别注意:已 await 完成的 task,cancel() 无效;只有 pending 状态的 task 才响应 cancel
拓扑调度真正的难点不在排序,而在异常路径下状态的一致性 —— 一个节点失败,整条链是否干净退出,取决于你有没有在每个 await 点做防御性检查。
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!










