← 返回题目列表

asyncio.Queue 如何实现异步生产者消费者模型?

高频 中等 第 6 / 27 题 更新于 2026/07/27
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 做收尾,配合哨兵值做优雅退出。