拒绝内存泄漏与连接打满:FastAPI异步SSE代理的生产级避坑实录

2次阅读
没有评论

共计 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 单次请求,它工作得完美无瑕。但只要放到真实公网或前端高并发场景,这段代码会引发三连击致命问题:

  1. 临时 Client 滥用与端口枯竭: 每个请求都 httpx.AsyncClient() 一次,底层完全无法复用 TCP/HTTP2 连接,在高并发短连接下,宿主机很快就会出现成千上万的 TIME_WAIT,直接报 Cannot assign requested address。
  2. 客户端主动断连时,后端请求继续跑(幽灵推理): 用户在前端看到回答不对劲,点了“停止生成”或直接关掉网页。FastAPI 默认的 StreamingResponse 并不会立刻杀死 event_generator 协程。你的网关依然在跟推理引擎通信,vLLM 还在持续生成 Token 浪费算力。
  3. 未捕获断连导致的协程悬挂与资源泄漏: 当生成器在 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。把这三处做扎实,你才能真正拥有一个稳如泰山的大模型转发层。

正文完
 0
评论(没有评论)