登录
推荐 文章 Go 技术 课程 下载 专题 AI
首页 >  文章 >  python教程

Python asyncio.Queue 如何实现有界生产者消费者

来源:17golang原创

时间:2026-09-07 02:09:17 330浏览 收藏

Python 的 asyncio.Queue 适合把生产者和消费者解耦,但真正能用于服务代码的关键不是“把数据放进去”,而是同时管住堆积上限、完成计数和退出信号。最小方案是使用 asyncio.Queue(maxsize=3),让生产者在队列满时等待;消费者每次 get() 后,无论处理成功还是失败,都在 finally 中调用 task_done();生产结束后调用 shutdown(),让消费者在清空队列后收到 QueueShutDown

要点速览
  • maxsize=0 表示无界,正整数才会对 put() 形成回压。
  • join() 等的是未完成任务数归零,不是简单等待队列变空。
  • shutdown() 适合正常收尾;immediate=True 会破坏通常的完成保证,只用于紧急终止。

maxsize 先把生产速度限制在可承受范围

Queue(maxsize=0) 是默认配置,队列不会因为达到容量而让 put() 等待。要让突发流量在内存里有明确边界,应设置正整数。队列满时,await queue.put(item) 会暂停当前生产协程,消费者取走数据后才继续,这就是应用层回压。

成员它解决的问题容易误解的边界
maxsize限制队列中允许保留的项目数不限制单个项目的大小
qsize()查看当前数量不能代替并发控制
put_nowait()满时立即失败需要处理 QueueFull
Python asyncio.Queue 的生产者、maxsize 队列容量与消费者之间的静态回压边界关系图
图1:查看生产者、maxsize、队列和消费者的静态边界,理解队列满时为何需要等待。

因此,maxsize 应按消费者并发数、单项内存占用和可接受等待时间共同估算。它是容量护栏,不是吞吐量承诺;若消费者被下游 I/O 拖慢,生产者仍会等待。

join 与 task_done 共同定义完成边界

join() 等待的是“所有已入队项目都完成处理”。每次 put() 会增加未完成计数,消费者取得项目后必须在处理结束时调用一次 task_done()。只要少调用一次,join() 就可能一直挂住;多调用一次则会抛出 ValueError

import asyncio

async def consumer(queue: asyncio.Queue):
    while True:
        try:
            item = await queue.get()
        except asyncio.QueueShutDown:
            # 队列已关闭且没有剩余项目,消费者可以退出
            break
        try:
            await handle(item)
        finally:
            # 每个 get 必须且只能配对一次 task_done
            queue.task_done()

async def producer(queue: asyncio.Queue, items: list[str]):
    for item in items:
        # 队列满时在这里等待,避免无界堆积
        await queue.put(item)

async def handle(item: str):
    await asyncio.sleep(0.01)

async def main():
    queue = asyncio.Queue(maxsize=3)
    worker = asyncio.create_task(consumer(queue))
    await producer(queue, ["a", "b", "c"])
    await queue.join()
    queue.shutdown()
    await worker

asyncio.run(main())

这段代码的完成边界很清楚:join() 返回只能说明每个项目都执行过对应的 task_done(),不能说明业务一定成功。生产环境中应在 handle() 周围记录异常或把失败项交给独立的重试策略。

shutdown 负责让消费者有序退出

Python 3.13 新增 shutdown()。默认的正常关闭会阻止新项目进入队列,但允许消费者继续取完已经入队的项目;队列清空后,新的 get() 会抛出 QueueShutDown。这比塞入特殊哨兵值更适合多个消费者,因为关闭状态由队列统一管理。

shutdown(immediate=True) 会立即排空队列并唤醒等待者,可能让 join() 在任务尚未真正处理时返回。它适合进程即将被强制终止的场景,不适合普通发布、定时任务或可恢复的消费流程。Python 3.12 及更早版本没有这个方法,需要继续使用哨兵值、取消任务或外围事件协调退出。

Python asyncio.Queue 中 producer、consumer、task_done、join、shutdown 与 QueueShutDown 的静态完成边界关系图
图2:对照生产、消费、完成计数和关闭信号的关系,区分正常收尾与立即终止。

常见问题

队列为空就代表所有任务完成了吗?

不一定。项目可能已经被 get() 取走但仍在处理,判断整体完成应等待 join()

为什么消费者要把 task_done 放进 finally?

处理函数抛异常时也要归还完成计数,否则生产者等待的 join() 没有机会结束;失败记录和重试应另行处理。

没有 Python 3.13 能调用 shutdown 吗?

不能直接调用。请使用版本兼容的哨兵值或任务取消方案,并把退出协议写进消费者协程。

声明:本文转载于:17golang原创 如有侵犯,请联系study_golang@163.com删除
相关阅读
更多>
最新阅读
更多>
课程推荐
更多>