共计 5631 个字符,预计需要花费 15 分钟才能阅读完成。
上个月我把自建的 vLLM 和几个商用 API 聚合到了一个统一的网关里,原本以为流式输出(Server-Sent Events, SSE)不过是 StreamingResponse 加个异步生成器就能搞定的事情。结果刚放给内部业务线使用,监控就直接报了警:显卡推理队列堆死、服务器内存持续攀升、上游网关连接数爆满。排查后发现三个致命问题:客户端断开后后台仍在盲跑推理浪费 Token、异步缓冲池缺乏背压机制导致单节点内存撑爆,以及不同客户端对 SSE 协议规范容忍度极差导致的解析崩溃。
今天这篇文章,我把这次重构出来的 LLM 流式代理网关方案完整脱敏整理出来。不谈虚的理论,直接切入核心工程痛点,带你搞定真正能在生产环境抗打的 Python 流式架构。
一、最简 Demo 为何会在生产环境“翻车”?
很多开发者在写流式转发时,直接照搬文档写出类似下面这样的代码:
# 常见但致命的简单写法
@app.post("/v1/chat/completions")
async def chat_proxy(payload: dict):
async def event_generator():
client = httpx.AsyncClient()
async with client.stream("POST", upstream_url, json=payload) as resp:
async select in resp.aiter_lines():
yield f"data: {select}\n\n"
return StreamingResponse(event_generator(), media_type="text/event-stream")
这套逻辑在单用户测试时毫无破绽,但一旦上并发,有三个巨坑必然踩中:
- 幽灵请求与算力浪费:用户在前端点击了“停止生成”或直接关掉了网页,客户端 TCP 链接已经 Reset(RST),但 FastAPI 后台的
aiter_lines()依然在欢快地拉取数据。上游昂贵的 GPU 资源或每千 Token 计费的商用 API 仍在疯狂空转。 - 背压(Backpressure)缺失引发 OOM:如果上游模型吐字极快,而前端网络恶劣(例如弱网移动端),TCP 接收窗口被填满,
StreamingResponse内部缓冲区会不受限制地暴涨,直至将网关内存耗尽。 - HTTPX 连接池耗尽:每次请求临时实例化
AsyncClient极度消耗文件描述符与连接握手开销;而如果复用全局 Client 但未精确控制连接生命周期,半开连接会迅速吃光池子。
二、架构重构:支持客户端断开感知的全异步代理
要解决上述问题,我们需要引入 Request.is_disconnected() 轮询检测、规范的 SSE 打包器、以及基于全局单例连接池的异步流水线。
1. 高性能流式网关核心实现
先来看我们重构后的核心中间层代码:
import json
import logging
from typing import AsyncGenerator
import httpx
from fastapi import FastAPI, Request, HTTPException
from fastapi.responses import StreamingResponse
logger = logging.getLogger("llm_proxy")
app = FastAPI()
# 共享且经过优化的上游 HTTP 连接池
UPSTREAM_CLIENT: httpx.AsyncClient = None
@app.on_event("startup")
async def startup_event():
global UPSTREAM_CLIENT
# 配置长连接、Keep-Alive 以及合适的连接池大小
limits = httpx.Limits(max_keepalive_connections=200, max_connections=500)
# LLM 推理首字延迟波动大,适当放宽 read 延迟
timeout = httpx.Timeout(connect=5.0, read=60.0, write=10.0, pool=10.0)
UPSTREAM_CLIENT = httpx.AsyncClient(limits=limits, timeout=timeout)
@app.on_event("shutdown")
async def shutdown_event():
await UPSTREAM_CLIENT.aclose()
async def safe_stream_generator(
request: Request,
upstream_url: str,
payload: dict
) -> AsyncGenerator[str, None]:
"""具备背压控制与断开感知的流式生成器"""
headers = {
"Authorization": "Bearer sk-upstream-secret-key",
"Content-Type": "application/json"
}
req = UPSTREAM_CLIENT.build_request("POST", upstream_url, json=payload, headers=headers)
try:
response = await UPSTREAM_CLIENT.send(req, stream=True)
if response.status_code != 200:
err_body = await response.aread()
logger.error(f"上游错误: {response.status_code} - {err_body.decode('utf-8')}")
yield f"data: {json.dumps({'error':'Upstream LLM error'})}\n\n"
return
async for chunk in response.aiter_raw():
# 关键:检查客户端是否已经断开
if await request.is_disconnected():
logger.warning("检测到客户端主动断开,立刻掐断上游请求避免算力浪费")
break
# 直接透传原装 SSE Chunk,减少反序列化 CPU 损耗
yield chunk
except httpx.ReadTimeout:
logger.error("上游推理超时")
yield f"data: {json.dumps({'error':'LLM timeout'})}\n\n"
except Exception as e:
logger.error(f"流式管道异常: {str(e)}")
yield f"data: {json.dumps({'error':'Internal proxy error'})}\n\n"
finally:
# 强制关闭响应流,通知上游中断当前未完成的推理
await response.aclose()
@app.post("/v1/chat/completions")
async def chat_completions(request: Request):
try:
payload = await request.json()
except Exception:
raise HTTPException(status_code=400, detail="Invalid JSON")
# 规范化响应头,阻止代理缓存与响应缓冲
headers = {
"Content-Type": "text/event-stream; charset=utf-8",
"Cache-Control": "no-cache, no-transform",
"Connection": "keep-alive",
"X-Accel-Buffering": "no", # 核心:关闭 Nginx 反向代理的 response buffering
}
upstream_url = "http://127.0.0.1:8000/v1/chat/completions" # 如本地 vLLM
return StreamingResponse(safe_stream_generator(request, upstream_url, payload),
headers=headers
)
2. 避坑重点:为什么一定要加 X-Accel-Buffering: no?
在生产架构中,FastAPI 外层几乎必挂 Nginx、Traefik 或 Caddy。我在第一次部署时踩过的最大坑就是:前端死活看不到“打字机”效果,而是傻等 10 秒钟后,整个回答突然“嘭”地一下全弹出来。
排查发现,Nginx 默认启用了 proxy_buffering,它把小数据包拼成了 4KB/8KB 的大包才一次性吐出。添加响应头 X-Accel-Buffering: no 可以直接告知 Nginx 关掉该路由的输出缓冲。如果用的是 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;
proxy_read_timeout 300s;
chunked_transfer_encoding on;
}
三、动态并发限流与令牌桶实战
不同于传统 Web 接口按 QPS 限流,LLM 流式请求的本质是 长连接 + 重资源占用。如果同时涌入 50 个流式请求,瞬间就能把显卡显存挤爆或让每秒吞吐断崖下跌。
我们采用 asyncio.Semaphore 配合基于 Redis 的动态并发插槽控制。下面是在代理层实现“最大并发生成数”的自研轻量级方案:
import asyncio
from contextlib import asynccontextmanager
class StreamConcurrencyManager:
def __init__(self, max_concurrent: int):
self.semaphore = asyncio.Semaphore(max_concurrent)
@asynccontextmanager
async def acquire_slot(self, timeout: float = 2.0):
try:
# 快速失败机制:若在限定时间内排不上号,直接拒绝,避免队列无限堆积
await asyncio.wait_for(self.semaphore.acquire(), timeout=timeout)
except asyncio.TimeoutError:
raise HTTPException(
status_code=429,
detail="推理集群过载,请稍后重试"
)
try:
yield
finally:
self.semaphore.release()
# 限制单实例最多同时处理 30 路流式推理
stream_limiter = StreamConcurrencyManager(max_concurrent=30)
@app.post("/v1/chat/limit-stream")
async def chat_limited_proxy(request: Request):
payload = await request.json()
async def wrapped_stream():
# 在生成器生命周期内持有信号量,直到客户端读取完毕或断开连接
async with stream_limiter.acquire_slot(timeout=1.5):
async for chunk in safe_stream_generator(request, "http://upstream/v1/chat", payload):
yield chunk
return StreamingResponse(wrapped_stream(),
media_type="text/event-stream"
)
四、压测与实测数据对比
我们使用 Locust 模拟 100 个客户端持续发起请求,并在生成中途有 30% 概率主动切断连接(模拟用户关闭界面或重新提问)。分别测试了“简易版转发”与“断开感知 + 缓冲优化版”在单台 8C16G 节点上的资源消耗表现:
| 测试维度 | 未优化的简易版 | 断开感知与背压优化版 | 改善幅度 |
|---|---|---|---|
| 客户端断开后上游浪费 Token | 平均每请求白跑 450 Tokens | 接近 0 Token(20ms 内终止) | 节省约 99% 无效开销 |
| 网关常驻内存 (RSS) | 从 180MB 稳步爬升至 1.4GB | 稳定在 220MB 左右波动 | 内存占用下降 84% |
| 首字响应延迟 (TTFT) | 并发高时因队列堵塞高达 3.8s | 稳定控制在 480ms 左右 | 响应速度提升约 7.9 倍 |
实测数据非常明显:在流式生成中,及早感知连接终结并级联取消(Cascade Cancellation),对于降本增效的价值远远大于盲目扩容显卡节点。
总结与避坑心得
构建高可用的 Python LLM 代理服务,难点从来不在于“把数据打印出来”,而在于异常流控制与资源释放边界:
- 永远不要信任客户端连接:只要写异步流式,务必注入
request.is_disconnected()轮询。它不仅保护你的代理服务器,更是在保护后端几万块钱一张的推理卡。 - 严控中间件反向代理:排查延迟问题,先用
curl -N -X POST ...直连 FastAPI 节点。如果直连打字顺畅、过 Nginx 后变卡顿,99% 是代理缓冲和 Keep-Alive 设置的问题。 - 全局连接池配置:千万不要在流式生成器函数内新建
httpx.AsyncClient()。使用应用生命周期内的单例池,并适度放大read_timeout,避免大模型输出长思考链(COT)时连接被强制掐死。
流式架构是 LLM 落地工程的“最后一公里”。把这段代码作为你的脚手架,能帮你避开绝大多数生产环境下的性能雷区。