asyncio.Queue 如何实现异步生产者消费者模型?
简化版
asyncio.Queue 是协程之间传递任务的异步队列,await queue.put() 放入任务,await queue.get() 获取任务,队列空或满时不会阻塞线程,而是挂起当前协程。它适合异步生产者消费者、爬虫下载、批量任务流水线。
详细版
示例:
import asyncio
async def producer(queue):
for i in range(100):
await queue.put(i)
await queue.put(None)
async def consumer(queue):
while True:
item = await queue.get()
try:
if item is None:
break
await handle(item)
finally:
queue.task_done()
async def main():
queue = asyncio.Queue(maxsize=10)
await asyncio.gather(producer(queue), consumer(queue))
关键点:
asyncio.Queue用在协程之间,不是线程安全队列;maxsize可以形成背压;get()后处理完要调用task_done();await queue.join()可以等待所有任务处理完;- 多个消费者退出时通常需要放多个哨兵值。
面试回答重点:异步队列把协程之间的共享状态改成消息传递,避免自己写复杂同步逻辑。
完整版教学
一、为什么异步程序也需要队列
异步程序里同样有生产者消费者问题:
- 一个协程抓取 URL;
- 多个协程下载页面;
- 一个协程解析结果;
- 一个协程写入数据库。
如果所有协程直接共享列表,就要考虑并发访问、空队列等待、限流、退出信号等问题。asyncio.Queue 把这些协调封装起来。
二、put 和 get 为什么要 await
普通线程队列的 put() / get() 可能阻塞线程。异步队列的操作是可等待的:
await queue.put(item)
item = await queue.get()
如果队列满了,生产者协程会挂起,事件循环去运行其他协程;如果队列空了,消费者协程挂起,事件循环也不会被卡住。
这正是异步队列和同步队列的关键区别。
三、maxsize 是背压
如果生产速度远大于消费速度,队列无限增长会吃掉内存。设置:
queue = asyncio.Queue(maxsize=100)
队列满时,生产者会等待消费者处理一些任务后再继续。这就是背压:下游处理不过来时,上游自动慢下来。
四、task_done 和 join
消费者取出任务后,处理完成要调用:
queue.task_done()
主协程可以等待所有任务完成:
await queue.join()
注意每次 get() 要对应一次 task_done()。如果忘记调用,join() 可能永远等下去;如果多调用,会抛异常。
五、多消费者如何退出
常见做法是放哨兵对象:
STOP = object()
for _ in range(worker_count):
await queue.put(STOP)
每个消费者拿到 STOP 后退出。多个消费者就需要多个 STOP,否则只有一个消费者退出,其他消费者可能一直等。
六、asyncio.Queue 和 queue.Queue 不要混用
queue.Queue 是线程间通信队列,它的阻塞方法会阻塞线程。异步协程之间应该用 asyncio.Queue。
如果你在 async 函数里直接调用同步队列的阻塞 get(),事件循环可能被卡住。不同并发模型的工具要匹配使用。
七、常见误区与追问
| 能力 | asyncio.Queue 表现 | 面试关注点 |
|---|---|---|
| 异步等待 | put/get 需要 await | 等待时让出事件循环 |
| 背压 | maxsize 限制积压 | 生产过快时阻塞生产者 |
| 完成确认 | task_done/join | 等任务处理完,不只是被取走 |
| 退出控制 | 哨兵值或取消任务 | 多消费者要发足够退出信号 |
- 误区:asyncio.Queue 和 queue.Queue 可以随便混用。 前者用于同一事件循环里的协程通信,后者用于线程通信;阻塞语义不同。
- 误区:get 到任务后不用 task_done。 如果使用
join()等待完成,每个成功取出的任务都应匹配一次task_done()。 - 误区:一个哨兵值能停止所有消费者。 多个消费者都阻塞在
get()时,通常需要给每个消费者一个哨兵,或设计明确的取消策略。 - 追问:maxsize 为什么是背压? 队列满时生产者
await put()会暂停,消费跟不上时压力不会无限堆到内存里。 - 追问:asyncio.Queue 是线程安全的吗? 它主要面向事件循环内协程,不应当把它当跨线程同步队列使用。
- 追问:消费者处理异常怎么办? 应确保异常路径也能调用必要的清理逻辑,避免未完成任务计数不归零导致
join()永久等待。
记忆钩子:异步队列的关键不是“先进先出”,而是把协程之间的等待、背压和完成确认统一起来。
八、加强记忆
asyncio.Queue 是协程世界的任务传送带:生产者 await put,消费者 await get,满了让生产者等,空了让消费者等,但不阻塞事件循环。配合 maxsize 做背压,配合 task_done / join 做收尾,配合哨兵值做优雅退出。