共计 4324 个字符,预计需要花费 11 分钟才能阅读完成。
前阵子我把团队内部的几个本地大模型(跑在 vLLM 和 Ollama 上的推理集群)以及商用 API 统一收口到一个自研的 Python 网关下。需求其实非常简单:做一层中转代理,提供鉴权、计费、Prompt 注入和 SSE(Server-Sent Events)打字机流式输出。原本以为这是个百来行 FastAPI 代码就能搞定的小活,结果压测和上线后接连被真实流量教做人——并发稍微拉到 500 以上,网关就开始报 PoolTimeout,宿主机内存一路缓涨不降,排查后才发现 Python 处理异步流式 HTTP 代理时,藏着不少教科书上根本不提的暗礁。
一、最直观的写法,为什么在生产环境必死?
先看一段最经典、网上一搜大把的 FastAPI SSE 代理 demo:
# 别在生产环境这么写!反面教材示例
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
import httpx
app = FastAPI()
@app.post("/v1/chat/completions")
async def chat_proxy(payload: dict):
client = httpx.AsyncClient(timeout=60.0)
async def event_generator():
async with client.stream("POST", "http://backend-llm:8000/v1/chat/completions", json=payload) as resp:
async for chunk in resp.aiter_bytes():
yield chunk
await client.aclose()
return StreamingResponse(event_generator(), media_type="text/event-stream")
如果在本地用 curl 单次请求,它工作得完美无瑕。但只要放到真实公网或前端高并发场景,这段代码会引发三连击致命问题:
- 临时 Client 滥用与端口枯竭: 每个请求都
httpx.AsyncClient()一次,底层完全无法复用 TCP/HTTP2 连接,在高并发短连接下,宿主机很快就会出现成千上万的TIME_WAIT,直接报Cannot assign requested address。 - 客户端主动断连时,后端请求继续跑(幽灵推理): 用户在前端看到回答不对劲,点了“停止生成”或直接关掉网页。FastAPI 默认的
StreamingResponse并不会立刻杀死event_generator协程。你的网关依然在跟推理引擎通信,vLLM 还在持续生成 Token 浪费算力。 - 未捕获断连导致的协程悬挂与资源泄漏: 当生成器在
yield处因为客户端断开而抛出CancelledError时,如果资源释放逻辑没写在严谨的finally块中,客户端连接根本来不及aclose(),连接句柄永远残留在内存里。
二、生产级架构设计:断连感知、背压与连接复用
要搞定这个场景,核心是三个机制:全局单例连接池 、FastAPI 底层断连感知(Request.is_disconnected),以及 严格的生命周期兜底。
1. 全局 Client 与连接池深度调优
HTTP/1.1 与 HTTP/2 的长连接复用是性能基石。我们需要利用 FastAPI 的 lifespan 管理全局单个 httpx.AsyncClient,并对底层 limits 进行压测校准:
from contextlib import asynccontextmanager
from fastapi import FastAPI
import httpx
# 全局 HTTP 客户端
http_client: httpx.AsyncClient = None
@asynccontextmanager
async def lifespan(app: FastAPI):
global http_client
limits = httpx.Limits(
max_keepalive_connections=200, # 保持的长连接数
max_connections=1000, # 最大并发连接上限
keepalive_expiry=30.0 # 空闲连接存活时间
)
# LLM 推理首字延迟波动大,必须细化超时时间
timeout = httpx.Timeout(
connect=5.0, # 建连超时
read=60.0, # 单个 Chunk 读取超时(不是整个对话超时)write=10.0,
pool=5.0 # 从连接池拿连接超时,防止请求无限积压
)
http_client = httpx.AsyncClient(limits=limits, timeout=timeout, http2=True)
yield
await http_client.aclose()
app = FastAPI(lifespan=lifespan)
2. 监听 Client 断开事件,主动掐断上游
这是整个代理最关键的代码。我们需要将 request: Request 传入,并在遍历生成器时主动检测客户端是否已经挂断。如果断开,立刻 break 并通知上游推理服务中断:
import asyncio
from fastapi import Request, HTTPException
from fastapi.responses import StreamingResponse
@app.post("/v1/chat/completions")
async def chat_proxy(request: Request):
payload = await request.json()
async def stream_generator():
req = http_client.build_request(
"POST",
"http://backend-llm:8000/v1/chat/completions",
json=payload,
headers={"Authorization": request.headers.get("Authorization", "")}
)
try:
resp = await http_client.send(req, stream=True)
if resp.status_code != 200:
err_body = await resp.aread()
yield f"data: {{\"error\": \"Upstream status {resp.status_code}: {err_body.decode('utf-8')}\"}}\n\n".encode()
return
async for chunk in resp.aiter_bytes():
# 核心探针:检测前端是否已经切断连接
if await request.is_disconnected():
# 客户端已断开,主动退出循环,触发 finally 关流
break
yield chunk
except asyncio.CancelledError:
# 框架级别强制取消任务触发
pass
except httpx.ReadTimeout:
yield b"data: {\"error\": \"LLM upstream read timeout\"}\n\n"
finally:
# 务必在 finally 里关闭 response stream,将连接还给连接池
if 'resp' in locals():
await resp.aclose()
return StreamingResponse(stream_generator(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"X-Accel-Buffering": "no" # 关掉 Nginx 缓冲区,必须直出!}
)
三、必须注意的网络拓扑与 Nginx 缓冲区暗坑
如果你的架构是 Browser -> Nginx -> FastAPI -> vLLM,你大概率会遇到“本地 curl 一切正常,浏览器打字机却卡成一块一块”的情况。这是 Nginx 默认的 Proxy Buffering 在搞鬼。
除了在 Python Header 里吐出 X-Accel-Buffering: no 外,在网关前端的 Nginx 配置文件中,关于 SSE 的反向代理配置必须精细化锁定:
location /v1/chat/completions {
proxy_pass http://127.0.0.1:8080;
# SSE 必备:关闭 Nginx 内部代理缓冲
proxy_buffering off;
proxy_cache off;
# 保持长连接,走 HTTP/1.1
proxy_http_version 1.1;
proxy_set_header Connection "";
# 针对大模型长思考的读写超时,防止中间切断
proxy_read_timeout 300s;
proxy_send_timeout 300s;
# 传递真实客户端状态
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
}
四、真实压力测试表现与数据对比
我在本地 8 核 16G 的测试机上,使用 locust 对优化前后的方案进行了 600 并发的压测对比。模拟用户行为:70% 的用户正常接收全部 500 Token 输出,30% 的用户在接收到前 30 个 Token 时主动切断请求(模拟取消提问)。
| 指标参数 | 优化前方案(频繁创建 Client / 无断连感知) | 优化后方案(全局池化 + 快速中断) |
|---|---|---|
| 网关平均内存占用 | 持续攀升至 1.8GB,并伴随 GC 停顿 | 稳定在 210MB ~ 260MB 区间 |
| 上游推理显存吞吐损耗 | 极高(幽灵请求持续压榨 GPU 计算卡) | 极低(上游收到 close 信号即刻释放 KV Cache) |
| P99 延迟 (首 Token 吐出) | 1420ms(连接池耗尽等待排队) | 118ms(长连接复用就绪) |
| 异常错误率 (5xx) | 8.4% (大多为 PoolTimeout 和连接拒绝) | 0.00% |
总结与避坑心得
在 Python 异步 Web 生态中,大吞吐的流式代理与传统的 JSON 响应有着本质不同。传统接口“请求进、数据出、生命周期结束”,而流式接口把请求拉长成了一个跨越数秒甚至数十秒的有状态管道。
如果你正在用 FastAPI 封装各类 LLM 网关,务必记住三条底线:永远不要在请求周期内随意 new 异步 HTTP 客户端;永远在流式迭代中加入 request.is_disconnected() 探针;最外层的反向代理(Nginx 或 Envoy)必须严密关闭 Buffer。把这三处做扎实,你才能真正拥有一个稳如泰山的大模型转发层。