登录
推荐 文章 Go 技术 课程 下载 专题 AI
首页 >  科技周边 >  人工智能

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

来源:17golang原创

时间:2026-08-29 21:29:04 308浏览 收藏

本地视觉模型服务刚接入批量图片上传时,最先暴露的往往不是模型精度,而是任务进入速度远高于 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删除
相关阅读
更多>
最新阅读
更多>
课程推荐
更多>