共计 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===