实战踩坑:基于 Python + vLLM 搭建高并发 OpenAI 兼容流式 Agent 网关

2次阅读
没有评论

共计 5407 个字符,预计需要花费 14 分钟才能阅读完成。

前段时间我给公司自建的 Agent 系统做架构升级,把后端的几台双卡 RTX 4090 节点通过 FastAPI 统一封装成兼容 OpenAI 的流式接口。最初以为写几个 StreamingResponse 和 httpx.AsyncClient 就能交差,结果在压测并发刚拉到 120 的时候,服务直接雪崩:客户端疯狂断流、内存以每分钟 500MB 的速度泄露,vLLM 的 PagedAttention 显存池甚至出现了莫名其妙的请求积压挂起。

这并不是 vLLM 本身推理性能拉胯,而是中间层这层“看似简单”的 Python 异步代理在处理 Server-Sent Events (SSE) 协议、连接生命周期管理以及背压(Backpressure)机制上踩了深坑。今天我把这套线上跑顺了的优化方案完整抽离出来,给同样在折腾私有化 LLM 网关的同学抄作业。

一、为什么简单的 httpx + FastAPI 转发会直接被打爆?

很多人刚开始写大模型流式代理,习惯直接在 FastAPI 路由里调用第三方客户端,然后无脑 yield,典型代码长这样:

# 反面教材:极其脆弱的流式透传
@app.post("/v1/chat/completions")
async def chat_proxy(request: Request):
    payload = await request.json()
    async def stream():
        async with httpx.AsyncClient() as client: # 灾难 1:每次请求创建 Client
            async with client.stream("POST", VLLM_URL, json=payload) as r:
                async for chunk in r.aiter_text():
                    yield chunk
    return StreamingResponse(stream(), media_type="text/event-stream")

如果你的场景只是一两个人自己玩玩,这段代码确实能跑。但一旦接入多人 Agent 或者跑并发压测,它会瞬间暴露出三个致命缺陷:

  • 连接池震荡与 FD 耗尽:每个进来的请求都新建一次 AsyncClient,TCP 连接无法复用,高并发下 Linux 系统的 TIME_WAIT 迅速拉满,直接报 Too many open files。
  • 客户端断开无法通知后端(显存严重浪费):当用户在前端点了“停止生成”或直接关掉网页,FastAPI 的生成器虽然会触发取消,但如果没有显式向 vLLM 发送 abort 信号,vLLM 会继续在显存中拼命生成剩余的几千个 Token,白白霸占显卡算力。
  • 无背压控制导致的内存泄漏:如果下游客户端网络环境较差(比如手机端 weak network),消费速率慢于 vLLM 吐字速率,没有缓冲区上限的异步生成器会把大量 Token 堆积在 Python 进程的内存队列中,导致 OOM。

二、生产级架构设计:单例连接池与硬核反向背压

为了解决上述问题,我们需要在 Python 层做到三件事:全局 HTTP 连接池复用、客户端掉线感知与推理取消注入、以及基于 asyncio.Queue 的有限流控缓冲。以下是提炼出来的生产级网关核心代码。

# gateway.py
import asyncio
import json
import logging
from contextlib import asynccontextmanager
from fastapi import FastAPI, Request, HTTPException
from fastapi.responses import StreamingResponse
import httpx

logger = logging.getLogger("llm_gateway")

class LLMGatewayManager:
    def __init__(self, backend_url: str):
        self.backend_url = backend_url
        self.client: httpx.AsyncClient | None = None

    async def init_client(self):
        # 针对长连接 SSE 进行定制化的连接池配置
        limits = httpx.Limits(max_keepalive_connections=200, max_connections=500)
        timeout = httpx.Timeout(connect=5.0, read=60.0, write=5.0, pool=5.0)
        self.client = httpx.AsyncClient(limits=limits, timeout=timeout)

    async def close_client(self):
        if self.client:
            await self.client.aclose()

manager = LLMGatewayManager("http://127.0.0.1:8000/v1/chat/completions")

@asynccontextmanager
async def lifespan(app: FastAPI):
    await manager.init_client()
    yield
    await manager.close_client()

app = FastAPI(lifespan=lifespan)

@app.post("/v1/chat/completions")
async def chat_completions_proxy(request: Request):
    try:
        body = await request.json()
    except Exception:
        raise HTTPException(status_code=400, detail="Invalid JSON body")

    # 针对 SSE 流式请求进行拦截与代理
    is_stream = body.get("stream", False)
    if not is_stream:
        # 非流式直接转发,此处省略,重点看流式逻辑
        resp = await manager.client.post(manager.backend_url, json=body)
        return resp.json()

    async def stream_generator():
        # 设置缓冲区深度,防止慢客户端拖垮内存
        queue = asyncio.Queue(maxsize=32)
        task_cancelled = False

        async def fetch_upstream():
            nonlocal task_cancelled
            try:
                # 开启流式通道请求上游 vLLM
                async with manager.client.stream("POST", manager.backend_url, json=body) as response:
                    if response.status_code != 200:
                        err_content = await response.aread()
                        await queue.put(f"data: {json.dumps({'error': err_content.decode()})}\n\n")
                        return

                    async for line in response.aiter_lines():
                        if task_cancelled:
                            break
                        if line:
                            # 队列满时会自动挂起,对上游产生背压限制
                            await queue.put(f"{line}\n\n")
            except httpx.ReadTimeout:
                logger.error("vLLM 推理超时")
                await queue.put("data: [ERROR: Timeout]\n\n")
            except Exception as e:
                logger.error(f"上游读取异常: {str(e)}")
            finally:
                # 放入哨兵标记,代表生成结束
                await queue.put(None)

        # 启动抓取任务
        fetch_task = asyncio.create_task(fetch_upstream())

        try:
            while True:
                # 检查客户端是否已经断开连接
                if await request.is_disconnected():
                    logger.warning("客户端提前断开连接,立即中断上游拉取")
                    task_cancelled = True
                    fetch_task.cancel()
                    break

                # 消费缓冲队列中的数据
                try:
                    item = await asyncio.wait_for(queue.get(), timeout=0.1)
                except asyncio.TimeoutError:
                    continue

                if item is None:
                    break

                yield item.encode("utf-8")
                queue.task_done()
        except asyncio.CancelledError:
            task_cancelled = True
            fetch_task.cancel()
            raise

    return StreamingResponse(stream_generator(),
        media_type="text/event-stream",
        headers={
            "Cache-Control": "no-cache",
            "Connection": "keep-alive",
            "X-Accel-Buffering": "no" # 极其关键:防止 Nginx 默认缓存 SSE 块
        }
    )

三、踩坑记录:老极客排障手记

代码跑起来只是第一步,在生产环境中你必然会遇到以下三个诡异隐患:

1. Nginx 默认缓冲导致打字机效果失效

在 Python 外部挂反向代理(如 Nginx)时,前端调用常常会遇到“等了五秒突然一口气喷出整段文本”的尴尬情况。这是因为 Nginx 默认开启了 proxy_buffering,它会攒够一个 buffer chunk(通常 4k 或 8k)才给客户端刷一次。

解决办法:除了在 Python 的 HTTP Headers 加上 X-Accel-Buffering: no 之外,Nginx 配置中也必须显式关闭:

location /v1/chat/completions {
    proxy_pass http://127.0.0.1:8080;
    proxy_http_version 1.1;
    proxy_set_header Connection "";
    proxy_buffering off;
    proxy_cache off;
    chunked_transfer_encoding on;
}

2. 替换系统默认的 asyncio 事件循环为 uvloop

CPython 自带的 asyncio 事件循环在处理数十万个轻量级协程切换与管道读写时,CPU 占用率偏高。换成用 C 编写的 uvloop,网络吞吐量通常能带来 20%~30% 的直接提升。启动脚本务必加上:

# 启动入口推荐使用 gunicorn + uvicorn 动态 worker
gunicorn gateway:app \
  --workers 4 \
  --worker-class uvicorn.workers.UvicornWorker \
  --bind 0.0.0.0:8080 \
  --loop uvloop \
  --timeout 120

3. vLLM 自身的显存碎片与并发限制

网关撑住了,如果 vLLM 崩了同样白搭。我在单卡 A100/4090 部署时,vLLM 启动参数有几个极其关键的配置:

python3 -m vllm.entrypoints.openai.api_server \
  --model /data/models/Qwen2.5-7B-Instruct \
  --served-model-name qwen-7b \
  --gpu-memory-utilization 0.92 \
  --max-model-len 8192 \
  --max-num-seqs 256 \
  --enable-chunked-prefill \
  --trust-remote-code

重点看 --enable-chunked-prefill。在 Agent 场景中,Prompt 往往包含很长的上下文,一旦多个超长 Prompt 同时到达,显存瞬间就会被 prefill 阶段撑爆触发 OOM。开启分块预填充后,长 Prompt 会被拆解与正在生成(decode)的请求穿插处理,P99 延迟显著降低。

四、实测压测数据对比

使用 Locust 对这套网关进行了持续 10 分钟的 200 并发流式压测,对比直接用裸脚本透传和经过本方案调优后的表现:

指标项 原始裸写方案 背压 + 单例连接池调优方案
首字延迟 (TTFT – P95) 1,420 ms 430 ms
网关进程常驻内存 (RSS) 从 180MB 飙升至 2.4GB (泄露) 稳定在 240MB 左右
异常中断率 (FD/Timeout) 14.8% 报错 0.00% 报错
客户端主动关闭后 GPU 负载释放 延迟 10~30 秒(继续空转) 立即释放(< 50ms)

总结与避坑心得

做大模型的基础设施,很多人把全部注意力都放在了模型量化、算子优化(vLLM / TensorRT-LLM)上,却往往忽略了最前端的 Python 胶水层。网络 I/O 虽轻,但在流式长链接、极高频数据切片的冲击下,任何一点细微的资源管理漏洞都会被成倍放大。

搭建此类高并发网关,牢记三点:绝对不要在单次请求内实例化 HTTP 客户端 、 务必利用有限队列为流式输出加上反向背压 、 严格监听 is_disconnected() 及时切断算力消耗。把这三件事做扎实,你的 Python 代理层就能抗住极其暴力的业务流量。

===EOF===

正文完
 0
评论(没有评论)