Python asyncio.Queue 给本地模型推理加背压:为什么任务堆积会挤爆显存
本地视觉模型服务刚接入批量图片上传时,最先暴露的往往不是模型精度,而是任务进入速度远高于 GPU 推理速度:进程看起来还活着,队列却越堆越长,显存和内存一起往上走。把任务交给无界列表只能把问题推迟,真正的处理办法是让生产者在固定容量的 asyncio.Queue 前停下来,并让 worker 在一次推理完成后明确调用 task_done()。
对本地模型推理,先把队列容量设成一个可解释的小数值,例如
maxsize=2,再用inference_mode()和queue.join()验证“提交受限、任务完成、显存可观察”这三件事是否同时成立。
asyncio.Queue(maxsize=2)满载后会让await put()等待,背压从队列边界开始生效。- 推理 worker 用
torch.inference_mode()包住前向计算,避免把推理过程误当成训练图。 - 每次
get()必须在finally中配对task_done(),否则queue.join()可能一直不返回。 torch.cuda.memory_allocated()与memory_reserved()只能帮助观察,不等于“调用清理函数就能增加可用显存”。
先复现:上传速度为什么会把推理服务拖进堆积
假设一个请求协程每次收到图片就创建任务,再把任务对象追加到普通 Python 列表。只要上传端比模型快,列表就没有自然的刹车点;图片、预处理张量和元数据会在消费者拿到之前继续占用内存。问题的关键不是“消费者要不要加速”,而是入口是否知道下游已经满了。
pending = []
async def submit(image_path):
pending.append(image_path)
return {"accepted": True, "pending": len(pending)}
这个写法连“最多允许等待几张图”都没有表达。改成有界队列后,容量本身成为系统的一部分:maxsize=2 表示最多保留两个尚未被 worker 取走的项目,第三次 put() 会等待,而不是悄悄扩大缓存。

把背压放在 asyncio.Queue.put 这一条边界上
下面的示例只保留一条清晰链路:生产者读取图片路径,queue.put() 负责排队,worker 调用 run_inference(),完成后用 task_done() 归还一个未完成任务。为了让示例能独立核对,模型推理函数用一个占位的异步包装表示真实推理调用;实际项目中可以把它替换为同步模型的线程池或进程池封装。
import asyncio
import torch
queue = asyncio.Queue(maxsize=2)
async def run_inference(image_path: str) -> dict:
await asyncio.sleep(0.2) # 示例:替换成真实的预处理与模型调用
return {"image": image_path, "status": "ok"}
async def producer(image_paths: list[str]) -> None:
for image_path in image_paths:
await queue.put(image_path)
async def worker() -> None:
while True:
image_path = await queue.get()
try:
with torch.inference_mode():
result = await run_inference(image_path)
print(result["image"], result["status"])
finally:
queue.task_done()
async def main(image_paths: list[str]) -> None:
task = asyncio.create_task(worker())
await producer(image_paths)
await queue.join()
task.cancel()
asyncio.run(main(["a.jpg", "b.jpg", "c.jpg", "d.jpg"]))
这里有两个容易漏掉的点。第一,queue.put() 的等待发生在生产者路径上,调用方需要把它当作正常的流控行为,而不是异常。第二,task.cancel() 只在 queue.join() 返回后执行,保证已经取出的任务都有完成标记。
队列满时看什么证据,才能确认背压真的生效
调试时不要只盯着最终输出。给入队和出队分别记录 queue.qsize(),你应当看到队列在 2 附近来回波动,而不是单调增长。为了观察等待边界,可以把入队前后的时间差记录下来:
import time
async def put_with_trace(image_path: str) -> None:
started = time.perf_counter()
await queue.put(image_path)
waited_ms = (time.perf_counter() - started) * 1000
print({"event": "enqueued", "image": image_path,
"queue_size": queue.qsize(), "waited_ms": round(waited_ms, 1)})
当 worker 变慢时,后续项目的 waited_ms 应该增加,同时 queue_size 不超过 2。这说明生产者在等待可用槽位。若队列大小仍然持续增加,通常是实际代码绕过了 put(),或者把多个无界中间列表藏在预处理、上传回调或批处理器里。

inference_mode 和显存指标各自能证明什么
torch.inference_mode() 是推理上下文,不是队列控制器;它解决的是自动求导相关的额外开销和张量记录边界。队列仍然要靠 maxsize 限制,不能因为加了 inference mode 就让入口无限接收任务。
显存排查可以记录两个不同指标:
torch.cuda.memory_allocated()更接近当前张量实际占用的显存。torch.cuda.memory_reserved()反映 PyTorch 缓存分配器管理的总量,可能高于 allocated。
如果任务已经完成但 reserved 没有立刻下降,不要直接判定为泄漏。PyTorch 文档说明,torch.cuda.empty_cache() 释放的是未占用的缓存块,并不会增加 PyTorch 自己可使用的显存;先确认仍有张量引用,再判断是否需要处理缓存碎片或生命周期。
三个常见坑:看似限流,实际没有完成闭环
只设置 maxsize,却在别处继续堆列表
队列前有界,不代表上传框架的批次缓存、重试列表和预处理缓存也有界。把每一层的容量写进配置,并在日志里区分“已接收”“已入队”和“已完成”,才能找到真正增长的容器。
忘记 task_done,join 永远等不到
queue.join() 等的是未完成任务计数归零。worker 在 get() 后遇到异常也必须走 finally,否则服务可能已经没有活任务,主协程却还在等待。
把 empty_cache 当成显存泄漏修复按钮
先用对象引用、批次生命周期和 allocated/reserved 的变化定位问题,再决定是否处理缓存。无界队列造成的输入堆积,不能靠周期性清缓存来补救。
相关问题
maxsize=2 是固定最佳值吗?
不是。它只是一个容易观察的起点。应按单个任务的内存占用、模型处理时长和可接受等待时间压测,再选择能解释清楚的容量。
worker 可以同时启动多个吗?
可以,但本地模型通常要先确认显存和模型线程安全边界。多个 worker 会提高并行度,也会增加批次和中间张量同时存在的机会。
为什么 queue.join() 返回后还要取消 worker?
worker() 是无限循环,队列清空只代表当前任务完成,不代表循环会自行退出。取消它可以让事件循环干净收尾。
收尾检查
把本地模型流水线交给别人维护前,至少核对四项:入口是否只用 await queue.put(),队列是否有明确的 maxsize,每次 get() 是否在 finally 中调用 task_done(),以及日志是否同时记录队列大小和 allocated/reserved。四项都能在代码和运行日志里对上,背压才不是一个停留在配置文件里的词。
大模型工具调用参数为空怎么办:Go 用 json.RawMessage 区分缺失与 null
- 上一篇
- 大模型工具调用参数为空怎么办:Go 用 json.RawMessage 区分缺失与 null
- 下一篇
- Kubernetes 1.37 Metrics API 晋级稳定:kubectl top 与自动扩缩容如何评估
-
- 科技周边 · 人工智能 | 2小时前 | 数据迁移 · 人工智能 · rag · embedding · 向量检索 · 向量数据库 向量维度 OpenAI Embeddings dimensions 索引迁移
- 向量检索为什么要统一维度:OpenAI Embeddings 的 dimensions 参数与索引迁移
- 105浏览 收藏
-
- 科技周边 · 人工智能 | 3小时前 | 人工智能 · 工程实践 · 模型评测 · LLM-as-a-Judge · 评测集 · 位置偏差 大模型评测 LLM-as-a-Judge 盲评 成对比较
- 大模型评测为什么偏爱更长答案:位置偏差、盲评与成对比较
- 329浏览 收藏
-
- 科技周边 · 人工智能 | 7小时前 | 人工智能 · mcp · 安全边界 · 协议设计 · URL MCP form Elicitation capability negotiation
- MCP capability negotiation 怎么确认客户端支持 elicitation:form 与 url 的分支判断
- 221浏览 收藏
-
- 前端进阶之JavaScript设计模式
- 设计模式是开发人员在软件开发过程中面临一般问题时的解决方案,代表了最佳的实践。本课程的主打内容包括JS常见设计模式以及具体应用场景,打造一站式知识长龙服务,适合有JS基础的同学学习。
- 543次学习
-
- GO语言核心编程课程
- 本课程采用真实案例,全面具体可落地,从理论到实践,一步一步将GO核心编程技术、编程思想、底层实现融会贯通,使学习者贴近时代脉搏,做IT互联网时代的弄潮儿。
- 516次学习
-
- 简单聊聊mysql8与网络通信
- 如有问题加微信:Le-studyg;在课程中,我们将首先介绍MySQL8的新特性,包括性能优化、安全增强、新数据类型等,帮助学生快速熟悉MySQL8的最新功能。接着,我们将深入解析MySQL的网络通信机制,包括协议、连接管理、数据传输等,让
- 500次学习
-
- JavaScript正则表达式基础与实战
- 在任何一门编程语言中,正则表达式,都是一项重要的知识,它提供了高效的字符串匹配与捕获机制,可以极大的简化程序设计。
- 487次学习
-
- 从零制作响应式网站—Grid布局
- 本系列教程将展示从零制作一个假想的网络科技公司官网,分为导航,轮播,关于我们,成功案例,服务流程,团队介绍,数据部分,公司动态,底部信息等内容区块。网站整体采用CSSGrid布局,支持响应式,有流畅过渡和展现动画。
- 485次学习
-
- ljg-skills
- ljg-skills 是李继刚开源的 AI 技能与提示词集合,面向大模型使用者整理了一批可复用的 prompt、角色设定和任务技能模板,适合用于学习提示词设计、搭建个人 AI 工作流和沉淀团队常用智能体能力。
- 5432次使用
-
- MELO音乐
- MELO音乐是一站式AI视频与音乐制作助手,对标suno, udio的高品质体验。提供伴奏生成、原创写词、无损导出、哼唱识曲、混音变声等全套音频与短视频编辑工具。无论是流行Kpop、电音说唱、民谣古风、摇滚儿歌还是商用轻音乐,MELO为你免费谱曲,轻松做同款!
- 4915次使用
-
- UniScribe
- UniScribe 是一款 AI 音视频转文字与内容整理工具,支持上传音频、视频文件或粘贴 YouTube 链接,自动生成转写文本、摘要、思维导图和关键问题,并支持多格式导出,适合会议记录、课程学习、访谈整理和内容创作复盘。
- 4838次使用
-
- 剧云
- 剧云是专业中文剧本创作平台,安全稳定运行十余年,集成AI编剧、剧本医生审核、人物小传、剧情关系图、大纲编写、多人协作、Word导入导出、版权管控功能,数据安全防护,轻松高效创作剧本。
- 5102次使用
-
- 万象有声
- 万象有声,一个专为有声创作者打造的新一代智能有声内容创作平台。平台提供专业的智能拆章、智能画本编辑、AI配音、AI生成音效、后期制作、智能对轨、智能审听等有声创作全流程工具,可以帮助创作者高效、低成本创作出引人入胜的有声作品。立即体验,让有声书制作更简单!
- 5058次使用
-
- go zero微服务实战性能优化极致秒杀
- 2022-12-27 207浏览
-
- go格式“占位符”输入输出 类似python的input
- 2023-01-19 346浏览
-
- Golang如何调用Python代码详解
- 2023-01-07 235浏览
-
- Go pprof 排查慢接口:别只会看火焰图,先把问题问对
- 2026-06-01 101浏览
-
- Go JSON v2 实战:别急着替换 encoding/json,先搞懂这些变化
- 2026-06-01 437浏览

