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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 14:27:02