Skip to content

本地模型/GPU 推理并发:管理 batch、显存和等待队列,而不是盲目多开 worker

对应视频:EP15《本地模型/GPU 推理并发》
视频入口:B站 EP15 定时稿(计划公开:2026-08-16 20:23)。 对应合集:Python 并发实战:从基础模型到语音工程。 本文为视频的工程展开版,补充代码模板、排查链路、状态字段和检查清单。 最后核验:2026/08/09 08:26:33

这篇文章是视频的工程展开版。视频负责建立问题现场、工具选择和排查链路;文章补充代码模板、状态字段、指标口径和落地检查清单。

GPU 推理 batch 队列图

这张图把 GPU 推理拆成入口队列、batcher、模型执行和指标反馈;并发管理的是 batch、显存和等待,而不是线程数量。

这篇解决什么问题

本地模型推理的并发核心是调度 GPU,而不是复制模型或无限排队。

Batch 策略

batch_size 提高吞吐,max_wait_ms 守住延迟,显存水位决定上限。离线和在线策略不同。

排查链路

看 gpu_util、memory、queue_wait_ms、batch_size、P95 和 OOM。

示例代码

下方保留可直接阅读的示例代码。

问题、解决方式和工具

问题解决方式工具
GPU 利用率不稳,显存却接近满载多开 worker 不一定提高吞吐,可能只是复制模型和放大显存。显存满 / 请求排队 / 延迟上升
有界队列保护显存和延迟请求队列解决入口压力,batcher 解决 GPU 吞吐。maxsize / drop/slowdown / priority
显存决定 batch 上限,也决定是否能多模型共存显存监控解决 OOM 风险和模型加载成本判断。nvidia-smi / reserved / allocated
离线追吞吐,在线守延迟目标不同,队列和 batch 策略也不同。offline / online / deadline

工程补充:怎么把这篇用到项目里

读图方式

  • 本地模型通常应该少量常驻 worker 复用模型,而不是为每个请求复制一份模型。
  • batch_size 提高吞吐,max_wait_ms 控制等待,二者共同决定延迟和 GPU 利用率。
  • 队列长度和显存水位决定是否继续接收任务、降 batch 或拒绝请求。

排查路径

  • nvidia-smi、框架内存指标、queue_wait、infer_time、p95 延迟一起看。
  • GPU 利用率低但 queue_wait 高,可能是 batcher、前处理或锁在阻塞。
  • 显存接近上限时,不要加 worker,先看模型份数、batch_size 和缓存释放。

问题、证据和解决动作

问题现场先查什么解决动作
多开 worker 后 OOM显存、模型实例数复用模型,减少 worker,限制 batch
GPU 利用率低preprocess_time、batch_wait合并 batch,优化前处理
延迟突然变高queue_wait、batch_size、max_wait_ms按延迟预算调 batcher
吞吐高但体验差p95/p99、超时率设置最大排队时间和拒绝策略

落地边界

  • GPU 并发的目标是稳定吞吐和可控延迟,不是让所有请求同时进入模型。
  • 显存水位要作为硬边界,不能等 OOM 后再靠重启恢复。

实战检查清单

  • 入口是否有容量边界,而不是无限提交。
  • 任务是否有状态字段,失败是否能恢复。
  • 每个重试是否有上限、退避和幂等保护。
  • 是否记录吞吐、P95、错误率、队列长度和关键资源水位。
  • 排查结果是否能指向具体动作:降并发、拆阶段、调 batch、隔离坏文件或回退方案。

结尾总结

最后把这一集收回来。

本地模型并发先要承认 GPU 和显存是硬边界。

入口用有界队列保护系统,batcher 用 batch_size 和 max_wait_ms 平衡吞吐与延迟,显存水位决定能不能继续加 batch 或多开模型。

排查时看 gpu_util、memory、queue_wait、batch_size 和 P95,才能知道问题在入口、预处理、batch 还是显存。

配套示例代码

源码文件:examples/ep15/gpu_batch_queue.py

python
import asyncio
import time

async def run_model(batch):
    await asyncio.sleep(0.03 + 0.005 * len(batch))
    return [f'result:{x}' for x in batch]

async def batcher(queue, batch_size=8, max_wait_ms=25):
    while True:
        item = await queue.get()
        batch = [item]
        deadline = time.perf_counter() + max_wait_ms / 1000
        while len(batch) < batch_size:
            timeout = max(0, deadline - time.perf_counter())
            if timeout == 0:
                break
            try:
                batch.append(await asyncio.wait_for(queue.get(), timeout))
            except asyncio.TimeoutError:
                break
        print(await run_model(batch))
        for _ in batch:
            queue.task_done()

async def main():
    queue = asyncio.Queue(maxsize=64)
    worker = asyncio.create_task(batcher(queue))
    for i in range(20):
        await queue.put(i)
    await queue.join()
    worker.cancel()

asyncio.run(main())
别急,先让缓存热一下。