共计 5263 个字符,预计需要花费 14 分钟才能阅读完成。
最近我把生产环境一个多 Agent 调度集群的接入层做了重构。原本前端直接对接内部部署的 vLLM 和 Ollama 实例,看似简单,但在上了一次线上活动后,并发请求瞬间把两台双路 A100 的推理节点打崩了。排查日志发现:不仅上游显存爆掉,代理层还堆满了大量挂起的死连接;更离谱的是,大量用户早已关闭网页或取消提问,但后端显卡依然在哼哧哼哧地把完整的 4096 个 Token 吐完。
很多开发者用 Python 写大模型流式代理(SSE, Server-Sent Events)时,习惯直接写个简单的 StreamingResponse。但在百级并发以上的长连接场景下,这种“玩具代码”会引发三大致命问题:连接池耗尽、客户端断开却不释放 GPU 算力、以及内存无界缓冲(缺乏背压控制)。今天这篇文章,我把这次重构沉淀下来的流式网关核心实现与调优参数完整拆解出来。
一、长连接流式代理的三个隐形暗礁
处理普通的 REST API,每个 HTTP 请求从接收到返回不过几十毫秒;但对于 LLM 流式输出,一个包含长思考链路或复杂回答的请求往往持续 15 到 45 秒。这种长连接会直接暴露出 Python 异步运行时的几个弱点:
- 文件描述符(FD)与连接池枯竭:默认的
httpx.AsyncClient只有 100 个最大连接和 20 个 keep-alive 连接。当 150 个用户同时发起长对话,第 101 个请求就直接排队甚至抛出PoolTimeout。 - “僵尸推理”吞噬算力:前端用户点击“停止生成”或直接关掉浏览器标签页,如果代理层没有感知到
http.disconnect并主动向上游发出cancel,推理引擎就会持续把这个 Prompt 运算到max_tokens。 - 缺乏流式背压(Backpressure):下游客户端网络差(比如手机端 weak 4G),接收 Chunk 的速度远慢于 GPU 推理速度。如果代理服务盲目向上游拉流并在内存中缓冲,很快整个容器的 RSS 内存就会直接 OOM。
二、架构设计与核心网关代码
我们采用 FastAPI + httpx + AnyIO 构建轻量代理。核心逻辑包括:单例长连接池管理、基于生成器的零缓冲实时管道转发,以及关键的 断连事件监听与协同取消机制。
import asyncio
import logging
from typing import AsyncGenerator
from fastapi import FastAPI, Request, HTTPException
from fastapi.responses import StreamingResponse
import httpx
app = FastAPI(title="LLM Stream Gateway")
logger = logging.getLogger("gateway")
# 全局共享的异步 Client,严禁每个请求重新实例化
upstream_client: httpx.AsyncClient | None = None
UPSTREAM_BASE_URL = "http://10.0.0.12:8000" # vLLM 或 Ollama 地址
UPSTREAM_TIMEOUT = httpx.Timeout(connect=5.0, read=60.0, write=5.0, pool=10.0)
UPSTREAM_LIMITS = httpx.Limits(max_connections=1000, max_keepalive_connections=200)
@app.on_event("startup")
async def startup():
global upstream_client
upstream_client = httpx.AsyncClient(
base_url=UPSTREAM_BASE_URL,
timeout=UPSTREAM_TIMEOUT,
limits=UPSTREAM_LIMITS,
http2=True # 启用 HTTP/2 减少连接建立握手
)
@app.on_event("shutdown")
async def shutdown():
if upstream_client:
await upstream_client.aclose()
async def forward_stream(
request: Request,
payload: dict
) -> AsyncGenerator[bytes, None]:
"""
流式代理转发生成器
实时监听下游断开信号,保证显存算力不被浪费
"""headers = {"Content-Type":"application/json","Authorization": request.headers.get("Authorization","")
}
try:
# 使用 stream() 避免将整个 body 加载进内存
async with upstream_client.stream(
"POST",
"/v1/chat/completions",
json=payload,
headers=headers
) as upstream_response:
if upstream_response.status_code != 200:
error_body = await upstream_response.aread()
logger.error(f"Upstream error [{upstream_response.status_code}]: {error_body}")
yield f"data: {{\"error\": \"Upstream inference failed\"}}\n\n".encode("utf-8")
return
# 按块实时迭代,配合客户端断连检测
async for chunk in upstream_response.aiter_raw():
# 检查客户端是否已经断开
if await request.is_disconnected():
logger.warning("Downstream client disconnected. Terminating upstream stream.")
break
if chunk:
yield chunk
except httpx.ReadTimeout:
logger.error("Upstream inference timeout.")
yield b"data: {\"error\": \"Gateway read timeout from model engine\"}\n\n"
except asyncio.CancelledError:
logger.info("Task cancelled by framework.")
raise
except Exception as e:
logger.exception(f"Unexpected streaming proxy error: {str(e)}")
yield b"data: {\"error\": \"Internal gateway error\"}\n\n"
@app.post("/v1/chat/completions")
async def chat_proxy(request: Request):
try:
payload = await request.json()
except Exception:
raise HTTPException(status_code=400, detail="Invalid JSON body")
# 强制确保业务处于 stream 模式
if not payload.get("stream", False):
raise HTTPException(status_code=400, detail="This gateway only serves streaming requests.")
return StreamingResponse(forward_stream(request, payload),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache, no-transform",
"Connection": "keep-alive",
"X-Accel-Buffering": "no" # 核心:禁止 Nginx/ 反代缓冲,保证即刻推送
}
)
三、实战踩坑与硬核优化点
1. 警惕 Nginx 的 proxy_buffering 截流
在调试过程中,本地跑测试脚本一切正常,一旦挂到 Nginx 后面,首字延迟(TTFT,Time to First Token)瞬间从 300ms 飙升到了 8 秒。前端根本不是“打字机”效果,而是模型跑完后一大坨数据直接砸过来。
排查发现是 Nginx 默认启用了响应缓冲机制。除了在 Python 代码中抛出 X-Accel-Buffering: no 响应头,网关前端的 Nginx 配置文件中必须显式设置:
location /v1/chat/completions {
proxy_pass http://python_gateway_backend;
proxy_http_version 1.1;
proxy_set_header Connection "";
# 彻底关闭缓冲,关键!proxy_buffering off;
proxy_cache off;
chunked_transfer_encoding on;
# 超时时间根据最长生成耗时放宽
proxy_read_timeout 180s;
proxy_send_timeout 180s;
}
2. 客户端断开检测与上下文取消机制
FastAPI 的 request.is_disconnected() 在轮询检测时会稍微占用少量 I/O 操作。在上面的代码中,每次收到 Upstream 的 Chunk 都会判断一次断开状态。一旦检测为 True,直接 break 退出。退出的瞬间,async with upstream_client.stream(...) 上下文管理器会被关闭,httpx 会主动向推理引擎(如 vLLM)切断 TCP 连接。vLLM 内部捕获到底层套接字断开后,会立即终止该 Request-ID 的 KV Cache 计算,从而将宝贵的算力让给其他并发任务。
3. 生产部署并发模型配置(Gunicorn + UvicornWorker)
不要用简单的 uvicorn main:app --workers 4 启动,在长连接场景下经常会出现 Worker 进程意外退出而无从排查。推荐使用 Gunicorn 托管,配合严格的 Worker 超时及文件句柄调优:
# gunicorn_conf.py
import multiprocessing
bind = "0.0.0.0:8080"
# 长连接流式应用以 I/O 等待为主,每个核可以适当提高并发,无需设过多纯 Worker 浪费内存
workers = multiprocessing.cpu_count() * 2 + 1
worker_class = "uvicorn.workers.UvicornWorker"
# 针对流式 API,将 worker 超时拉长,避免长文本推理被当成超时进程杀死
timeout = 300
keepalive = 75
# 调整最大并发接受连接数
worker_connections = 4096
max_requests = 50000
max_requests_jitter = 5000
启动前别忘了把系统级别的最大打开文件数(ulimit)调大,不然高并发下会出现 OSError: [Errno 24] Too many open files:
ulimit -n 65535
gunicorn -c gunicorn_conf.py app:app
四、压测实测对比数据
我们在双路 64 核服务器上部署该代理网关,后端挂载 2 个 vLLM 实例(DeepSeek-V2-Lite),使用 Locust 进行并发长文本压测(每个请求预期生成约 800 Tokens):
| 压测模式 | 并发用户数 | 首字延迟 (TTFT P99) | 算力空转率 (取消请求) | 网关内存占用 (RSS) |
|---|---|---|---|---|
| 默认方案 (无池化 / 无断连中断) | 200 | 4820 ms | 100%(依然算完) | ~1.8 GB (内存溢出边缘) |
| 优化方案 (连接池复用 + 协同取消) | 200 | 430 ms | < 2%(断连即终止) | 稳定在 210 MB |
从实测数据可以非常直观地看出:全局单例连接池消除了 TCP 频繁握手以及内存中累积的大量一次性 Client 垃圾;而 断连协同中断 在大流量压测下,把由于客户端掉线、主动取消带来的无效显存占用降低了 90% 以上,直接让后端推理吞吐量翻倍。
总结与避坑心得
大模型工程落地中,网关层看似只是一层数据搬运工,但从短文本、低延迟的 Web2 请求跨越到长耗时、高占用、流式推送的 LLM 场景,底层的网络 I/O 逻辑已经彻底变了。
归结起来,构建生产级 Python LLM 网关只需守死四条军规:全局复用 httpx 连接池、绝不在反向代理层开 buffer、实时捕获客户端断连通知上游 abort、把 ulimit 和 worker_connections 拉满。按照这套架构跑下来,基本单机几百兆内存就能轻松稳住数千个长文本推理连接,再也不用担心下游一次网络抖动拖垮一整张显卡了。