Python asyncio中add_done_callback失效及异步限流实现问题
问题根因分析
- KeyError问题:本质是Python闭包的*迟绑定(late binding)*特性导致的。你在循环中定义的lambda回调捕获的是变量
t的引用,而非当前迭代的t的值,所有回调执行时读取的t都是循环结束时最后一个任务的引用,因此会出现删除错误任务、重复删除的情况,最终抛出KeyError。 - 残留任务问题:同样由上述闭包问题导致,本该被删除的旧任务没有被回调正确删除,一直残留在
task_set中,最终会导致len(task_set)永远大于限流阈值,程序卡住不再生成新任务。
修复后的代码实现
import asyncio def main(): async def task(acc_id): print(f"acc_id {acc_id} received the task") await asyncio.sleep(5) print(f"acc_id {acc_id} finished the task") async def loop_tasks(): task_set = set() limit = 10 for acc_id in range(1000000): # 达到限流阈值时等待任意一个任务完成,避免轮询空耗CPU if len(task_set) >= limit: done, pending = await asyncio.wait( task_set, return_when=asyncio.FIRST_COMPLETED ) # 批量移除已完成的任务 task_set.difference_update(done) t = asyncio.create_task(task(acc_id)) task_set.add(t) # 通过默认参数绑定当前迭代的任务对象,用discard避免极端情况的KeyError t.add_done_callback(lambda _, current_t=t: task_set.discard(current_t)) # 所有任务生成完成后,等待剩余任务全部执行完毕再退出 if task_set: await asyncio.wait(task_set) asyncio.run(loop_tasks()) if __name__ == "__main__": main()
核心修改说明
- 修复闭包迟绑定:给lambda添加默认参数
current_t=t,每次迭代时都会把当前的任务对象绑定到回调的参数中,避免所有回调共享同一个t引用。 - 异常安全删除:将
remove方法替换为discard,就算任务已经不在集合中也不会抛出KeyError。 - 优化限流等待逻辑:替换轮询
sleep(0)的方案,使用asyncio.wait的FIRST_COMPLETED模式,只有当有任务完成时才会唤醒继续生成新任务,大幅降低CPU空耗。 - 补充收尾逻辑:所有任务生成完成后,等待剩余在飞的任务全部执行完毕再退出,避免程序提前终止导致任务被强制取消。
如果你之前用Semaphore不符合预期,大概率是因为传统Semaphore用法需要提前创建所有任务再靠信号量限流,面对百万级任务时会占用极高内存。上述实现属于轻量化的生产者-消费者变体,一边生成任务一边限流,内存占用始终维持在极低水平,更适合大数量级的任务场景。
内容的提问来源于stack exchange,提问来源于stack exchange
相关产品推荐
相关产品推荐

