当前位置:首页 > 文章列表 > 科技周边 > 人工智能 > Python asyncio.Queue 给本地模型推理加背压:为什么任务堆积会挤爆显存

Python asyncio.Queue 给本地模型推理加背压:为什么任务堆积会挤爆显存

来源:17golang原创 2026-08-29 21:29:04 0浏览 收藏

本地视觉模型服务刚接入批量图片上传时,最先暴露的往往不是模型精度,而是任务进入速度远高于 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() 会等待,而不是悄悄扩大缓存。

Python asyncio.Queue 与本地模型推理的调用链:生产者进入有界队列后由模型推理 worker 取出处理

把背压放在 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(),或者把多个无界中间列表藏在预处理、上传回调或批处理器里。

队列满到等待再到 task_done 的状态变化:Python 本地模型推理用 qsize 和显存指标检查稳定性

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。四项都能在代码和运行日志里对上,背压才不是一个停留在配置文件里的词。

版本声明
本文转载于:17golang原创 如有侵犯,请联系study_golang@163.com删除
大模型工具调用参数为空怎么办:Go 用 json.RawMessage 区分缺失与 null大模型工具调用参数为空怎么办:Go 用 json.RawMessage 区分缺失与 null
上一篇
大模型工具调用参数为空怎么办:Go 用 json.RawMessage 区分缺失与 null
Kubernetes 1.37 Metrics API 晋级稳定:kubectl top 与自动扩缩容如何评估
下一篇
Kubernetes 1.37 Metrics API 晋级稳定:kubectl top 与自动扩缩容如何评估
查看更多
最新文章
查看更多
课程推荐
  • 前端进阶之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推荐
  • ljg-skills -
    ljg-skills
    ljg-skills 是李继刚开源的 AI 技能与提示词集合,面向大模型使用者整理了一批可复用的 prompt、角色设定和任务技能模板,适合用于学习提示词设计、搭建个人 AI 工作流和沉淀团队常用智能体能力。
    5432次使用
  • MELO音乐 - AI 音乐生成平台,支持多模态创作能力
    MELO音乐
    MELO音乐是一站式AI视频与音乐制作助手,对标suno, udio的高品质体验。提供伴奏生成、原创写词、无损导出、哼唱识曲、混音变声等全套音频与短视频编辑工具。无论是流行Kpop、电音说唱、民谣古风、摇滚儿歌还是商用轻音乐,MELO为你免费谱曲,轻松做同款!
    4915次使用
  • UniScribe - AI 免费在线音视频转文字平台
    UniScribe
    UniScribe 是一款 AI 音视频转文字与内容整理工具,支持上传音频、视频文件或粘贴 YouTube 链接,自动生成转写文本、摘要、思维导图和关键问题,并支持多格式导出,适合会议记录、课程学习、访谈整理和内容创作复盘。
    4838次使用
  • 剧云 - 免费 AI 智能中文剧本创作平台
    剧云
    剧云是专业中文剧本创作平台,安全稳定运行十余年,集成AI编剧、剧本医生审核、人物小传、剧情关系图、大纲编写、多人协作、Word导入导出、版权管控功能,数据安全防护,轻松高效创作剧本。
    5102次使用
  • 万象有声 - AI 一站式有声内容创作平台
    万象有声
    万象有声,一个专为有声创作者打造的新一代智能有声内容创作平台。平台提供专业的智能拆章、智能画本编辑、AI配音、AI生成音效、后期制作、智能对轨、智能审听等有声创作全流程工具,可以帮助创作者高效、低成本创作出引人入胜的有声作品。立即体验,让有声书制作更简单!
    5058次使用
微信登录更方便
  • 密码登录
  • 注册账号
登录即同意 用户协议隐私政策
返回登录
  • 重置密码