当前位置:首页 > 文章列表 > 文章 > python教程 > Python asyncio.Queue 如何实现有界生产者消费者

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

来源:17golang原创 2026-09-07 02:09:17 0浏览 收藏

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删除
Go select 里 default 为什么让循环占满 CPUGo select 里 default 为什么让循环占满 CPU
上一篇
Go select 里 default 为什么让循环占满 CPU
Go 怎么安全地批量重命名文件并支持失败回滚
下一篇
Go 怎么安全地批量重命名文件并支持失败回滚
查看更多
最新文章
查看更多
课程推荐
  • 前端进阶之JavaScript设计模式
    前端进阶之JavaScript设计模式
    设计模式是开发人员在软件开发过程中面临一般问题时的解决方案,代表了最佳的实践。本课程的主打内容包括JS常见设计模式以及具体应用场景,打造一站式知识长龙服务,适合有JS基础的同学学习。
    543次学习
  • GO语言核心编程课程
    GO语言核心编程课程
    本课程采用真实案例,全面具体可落地,从理论到实践,一步一步将GO核心编程技术、编程思想、底层实现融会贯通,使学习者贴近时代脉搏,做IT互联网时代的弄潮儿。
    516次学习
  • 简单聊聊mysql8与网络通信
    简单聊聊mysql8与网络通信
    如有问题加微信:Le-studyg;在课程中,我们将首先介绍MySQL8的新特性,包括性能优化、安全增强、新数据类型等,帮助学生快速熟悉MySQL8的最新功能。接着,我们将深入解析MySQL的网络通信机制,包括协议、连接管理、数据传输等,让
    500次学习
  • JavaScript正则表达式基础与实战
    JavaScript正则表达式基础与实战
    在任何一门编程语言中,正则表达式,都是一项重要的知识,它提供了高效的字符串匹配与捕获机制,可以极大的简化程序设计。
    487次学习
  • 从零制作响应式网站—Grid布局
    从零制作响应式网站—Grid布局
    本系列教程将展示从零制作一个假想的网络科技公司官网,分为导航,轮播,关于我们,成功案例,服务流程,团队介绍,数据部分,公司动态,底部信息等内容区块。网站整体采用CSSGrid布局,支持响应式,有流畅过渡和展现动画。
    485次学习
查看更多
AI推荐
  • SuperCLUE中文大模型评测基准:功能、能力维度与应用指南
    SuperCLUE
    SuperCLUE是权威的中文大语言模型综合评测基准,涵盖语言理解、知识应用、AI Agent智能体及安全性等12项核心能力。通过多轮对话与客观测试,定期发布榜单与技术报告,为模型研发、优化及行业选型提供科学依据。
    167次使用
  • C-Eval中文评测基准:大语言模型多学科能力评估指南
    C-Eval
    深入了解C-Eval中文评估套件,涵盖52个学科与4级难度。本文详解其功能特点、Zero-shot/Few-shot使用方法及代码示例,助您全面评测LLM中文理解与泛化能力。
    93次使用
  • AI Prompt Library:免费AI提示词库,助力ChatGPT高效创作与营销
    AI Prompt Library
    探索AI Prompt Library免费资源库,涵盖营销、写作及多场景AI提示词。兼容ChatGPT、Claude等工具,一键复制优化输出,提升工作效率。
    16次使用
  • LangGPT提示词框架:结构化Prompt设计方法与开源工具指南
    LangGPT
    LangGPT是一种受编程语言启发的结构化提示词设计工具,提供双层框架、模块化模板及变量功能,帮助用户高效编写高质量Prompt。该项目已在GitHub免费开源,适用于内容创作、编程辅助等多场景。
    29次使用
  • ClickPrompt:AI提示词生成与优化工具,支持Stable Diffusion、ChatGPT及代码辅助
    ClickPrompt
    ClickPrompt是一款专为AI提示词编写者设计的开源在线工具,支持Stable Diffusion绘图、ChatGPT对话及GitHub Copilot代码辅助。提供Prompt自动生成、一键运行、社区分享及可视化优化功能,帮助用户高效获取精准AI输出。
    61次使用
微信登录更方便
  • 密码登录
  • 注册账号
登录即同意 用户协议隐私政策
返回登录
  • 重置密码