共计 4664 个字符,预计需要花费 12 分钟才能阅读完成。
上个月我给团队的私有化大模型服务搭了一套代理网关,本以为只是简单的“收请求 -> 调 vLLM/Ollama -> 吐 SSE 流”,结果刚灰度推上线就被狠狠上了一课:压测并发刚拉到 300,Python 进程的事件循环就开始严重延迟,更有甚者,前端用户点取消生成或者直接关闭网页,后端 vLLM 还在傻乎乎地把整段 Prompt 推理完,直接把两张 A100 的显存与计算算力白白榨干。
如果你也在用 Python 做大模型网关、Agent 中转或流式接口(Server-Sent Events),普通的 requests 甚至默认配置的 StreamingResponse 绝对撑不住生产流量。今天我把这套经过实际业务压测打磨的流式网关核心实现、背压控制机制以及踩坑调优经验彻底拆开讲透。
一、为什么简单的流式转发会直接崩掉?
很多同学在写 FastAPI 流式代理时,网上抄的代码大概长这样:
# 典型反面教材:看似异步,实则暗藏隐患
@app.post("/v1/chat/completions")
async def chat_proxy(request: Request):
client = httpx.AsyncClient()
req = client.build_request("POST", UPSTREAM_URL, json=await request.json())
r = await client.send(req, stream=True)
return StreamingResponse(r.aiter_raw(), media_type="text/event-stream")
这种写法在线上跑起来,主要会踩三个致命大坑:
- 客户端断连无法感知(显存幽灵):用户关闭标签页后,FastAPI 默认不会主动打断下游的异步生成器,网关依然在吃上游推理引擎的流,上游完全不知道客户端已断开,GPU 继续满负荷吐字,极度浪费计算资源。
- 内存膨胀与背压缺失(Backpressure Failure):大模型首字出来后吐字飞快,但如果客户端处于弱网环境、TCP 窗口打满,中间网关的 Python 进程为了缓存未下发的数据块,内存会像气球一样瞬间鼓起来。
- 连接池耗尽与事件循环阻塞:临时创建
AsyncClient会让连接频繁建立与销毁;如果下游用了同步解析库或不当的async for循环逻辑,极易引发事件循环延迟(Event Loop Lag)。
二、高并发流式代理网关的工程级实现
要解决上述痛点,网关必须具备三个能力:复用长连接池 、 监听下游断连并快速通知上游中断(Cancel Propagation)、控制数据块滑动窗口与背压。
下面是我在线上环境验证过的精简版核心代理代码:
import asyncio
import logging
from contextlib import asynccontextmanager
from fastapi import FastAPI, Request, HTTPException
from fastapi.responses import StreamingResponse
import httpx
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("stream-gateway")
UPSTREAM_LLM_URL = "http://10.0.0.12:8000/v1/chat/completions"
# 全局复用长连接池,控制并发上限
class HttpClientHolder:
client: httpx.AsyncClient = None
client_holder = HttpClientHolder()
@asynccontextmanager
async def lifespan(app: FastAPI):
# 根据上游并发承载能力设定连接池参数
limits = httpx.Limits(
max_connections=1000,
max_keepalive_connections=200,
keepalive_expiry=30.0
)
# 大模型长文本生成可能耗时极长,合理设置超时
timeout = httpx.Timeout(connect=5.0, read=120.0, write=5.0, pool=10.0)
client_holder.client = httpx.AsyncClient(limits=limits, timeout=timeout)
yield
await client_holder.client.aclose()
app = FastAPI(lifespan=lifespan)
async def stream_generator(request: Request, payload: dict):
client = client_holder.client
req = client.build_request("POST", UPSTREAM_LLM_URL, json=payload)
try:
response = await client.send(req, stream=True)
if response.status_code != 200:
err_body = await response.aread()
logger.error(f"上游返回异常状态码: {response.status_code}, 内容: {err_body}")
yield f"data: {{\"error\": \"Upstream returned {response.status_code}\"}}\n\n".encode("utf-8")
return
async with response:
async for chunk in response.aiter_raw():
# 关键一步:主动探测下游客户端是否断开
if await request.is_disconnected():
logger.warning("检测到客户端断开连接,立即切断与上游连接!")
break
if chunk:
yield chunk
except httpx.RequestError as exc:
logger.error(f"调用上游服务通信失败: {str(exc)}")
yield f"data: {{\"error\": \"Gateway Connection Error\"}}\n\n".encode("utf-8")
except asyncio.CancelledError:
logger.info("当前任务已被取消,释放相关上下文")
raise
finally:
logger.debug("流式代理周期结束")
@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 payload")
# 强制将 stream 置为 True 走流式转发
payload["stream"] = True
return StreamingResponse(stream_generator(request, payload),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"X-Accel-Buffering": "no" # 关键:避免 Nginx 默认缓存 SSE 流
}
)
三、核心调优细节与生产避坑
上面这段代码看似简短,但包含了三个极为关键的实操细节:
1. request.is_disconnected() 的触发频率
在 async for chunk in response.aiter_raw() 每次循环迭代中,我们都检查了一次连接状态。大模型返回是以 token 为单位的碎片数据,频率通常在几十到几百毫秒之间,这个检测间隔刚好能在用户关闭网页的 100ms 内中断下游循环。只要跳出 async with response 上下文,httpx 会自动关闭底层 TCP 连接,上游(如 vLLM)检测到客户端断开就会立刻调用自身 cancellation handler 停止 decoding,显存算力当场释放。
2. Nginx 反向代理的“吸血”缓冲坑
我在测试阶段发现,客户端根本拿不到流式输出,必须等模型全部输出完才一次性刷出来。排查了一圈网关代码都没问题,最后在反代层发现了罪魁祸首:Nginx 默认开启了 proxy_buffering。必须在响应头中带上 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 "";
# 彻底关掉缓存,保证实时 SSE 吐字
proxy_buffering off;
proxy_cache off;
chunked_transfer_encoding on;
# 调大长连接保活和响应超时
proxy_read_timeout 600s;
proxy_send_timeout 600s;
}
3. 避免单进程事件循环被同步代码“锁死”
如果在流式迭代器里加入日志记录、安全合规敏感词检测、Prompt 审计,千万不要直接调用同步库(如 time.sleep 或本地重度同步加解密函数)。可以使用 asyncio.to_thread 将计算密集型操作下放至独立线程池,否则只要单次卡住 50ms,几十路流的吐字帧率就会像幻灯片一样卡顿。
四、压测数据实测与指标表现
我们采用 Locust 对单节点网关(4 Core 8G,配置 2 个 Uvicorn Worker)进行压测,后端接入固定延时 Mock 的 SSE 生成源,测试在不同并发场景下的表现:
| 方案架构 | 并发连接数 | 首字延迟 (TTFT P95) | 断连即时释放率 | 网关内存占用 |
|---|---|---|---|---|
| 普通无状态转发(无复用池) | 200 | 1450ms | 0%(完全不释放) | ~1.8 GB(持续泄露) |
| 全局 httpx 连接池 + 状态检测 | 200 | 180ms | 99.2% | ~220 MB |
| 全局连接池 + 缓冲关闭优化 | 500 | 215ms | 99.8% | ~310 MB |
从实测数据可以非常清晰地看到:一旦加入了长连接池管理与断连状态拦截,TTFT(Time To First Token,首字输出时间)从将近 1.5 秒直接压进 200 毫秒以内,更关键的是内存占用和计算资源的及时止损,完全不在一个量级。
总结与避坑心得
大模型时代的代理网关,已经从传统的“短报文高 QPS 吞吐”演化成了“超长连接、流式传输、双向状态同步”的全新形态。如果你正准备用 Python 重构网关服务,记住这三条铁律:
- 永远不要在每次请求内部随手
httpx.AsyncClient(),必须在应用生命周期全局维护连接池。 - 必须在流生成器内部主动轮询
request.is_disconnected(),把取消信号传递给上游,这是保住 GPU 显存的核心手段。 - 全链路关闭数据缓冲,包括网关自身的响应配置和前端各类 Nginx/Traefik 负载均衡器。
用这套骨架搭建的网关,代码量少且高度可控,不用引入 Java/Go 那些重型网关框架,小团队单台服务器即可轻松扛住数千路实时的大模型长连接交互。