用 asyncio.Queue 搭配常驻协程消费者可实现动态添加任务的异步工作队列,核心是解耦生产与消费、避免阻塞事件循环,并通过 task_done() 和 join() 确保优雅关闭。

用 asyncio.Queue 搭配协程消费者,就能轻松实现支持动态添加任务的异步工作队列。核心是让生产者和消费者解耦,且队列本身不阻塞事件循环。
用 asyncio.Queue 作为任务中转站
asyncio.Queue 是线程安全、协程友好的队列,天然适配异步环境。它支持 put_nowait()/put() 和 get()/get_nowait(),能自动处理协程等待逻辑。
- 创建队列时可设 maxsize(如 asyncio.Queue(maxsize=10)),超限时 put() 会挂起,避免内存失控
- 调用 await queue.put(task) 添加任务,await queue.get() 获取任务,获取后需手动调用 task_done()
- 多个消费者可并发 await queue.get(),Queue 内部自动负载均衡
启动常驻消费者协程
消费者应作为长期运行的协程存在,持续从队列取任务并执行,而不是为每个任务单独启一个协程(易失控)。
- 用 asyncio.create_task() 启动一个或多个消费者协程,例如:asyncio.create_task(worker(queue))
- 消费者内部用 while True + await queue.get() 循环监听,捕获异常避免崩溃退出
- 处理完任务后必须调用 await queue.task_done(),否则 join() 无法判断完成状态
动态添加任务无需重启队列
只要队列对象还活着,任何地方(包括其他协程、HTTP 请求处理函数、定时器回调)都能随时 await queue.put(task)。
- 任务可以是普通函数、协程对象,或封装了参数和上下文的可调用对象(如 dataclass 或 dict)
- 例如:await queue.put({"func": send_email, "args": ("user@example.com", "Hi!")})
- 消费者拿到后用 await func(*args) 或直接 await task 执行(若 task 是协程)
优雅关闭与等待完成
停止生产后,需等待所有任务处理完毕再退出,避免任务丢失。
- 调用 await queue.join() 阻塞直到所有已入队任务都被 task_done() 标记完成
- 可配合 asyncio.Event 或标志位通知消费者退出循环(例如 put sentinel 值或检查全局 stop_flag)
- 实际关闭时,先停止生产,再 join,最后 cancel 消费者任务(如有必要)

















