FastAPI+httpx 流式大模型网关实战:避开内存泄露与连接池耗尽深坑

9次阅读
没有评论

共计 4164 个字符,预计需要花费 11 分钟才能阅读完成。

给业务系统做大模型统一接入层时,绝大多数同学的第一反应都是:用 FastAPI 挂一个反向代理,把前端的请求转发给内网的 vLLM、Ollama 或者外部 OpenAI 兼容接口,再用 StreamingResponse 把 SSE(Server-Sent Events)流实时吐回去。

这套逻辑写个 demo 只需要二十行代码。然而,上个月我们把一套看似正常的网关推上压测环境,模拟 200 个并发长文本生成时,服务直接挂掉:容器内存从 80MB 直线飙破 3.5GB 触发 OOM,终端刷屏 httpx.PoolTimeout,后台监控里 GPU 的利用率居高不下,即使客户端早已断开连接。

这篇文章记录我在生产环境中排查并解决这几个致命问题的真实过程,顺便给出可以直接压测跑通的生产级网关代码模板。

第一坑:httpx.AsyncClient 的连接池与超时默认值

很多新手写 FastAPI 转发请求时,习惯在路由函数内直接 async with httpx.AsyncClient() as client:。这在低频短连接场景下问题不大,但在高并发流式场景下是灾难——每次请求都在反复创建 TCP 和 TLS 握手,连接无法复用。

更致命的是,有人虽然用了全局单例 Client,但没改默认参数:

# 危险的默认配置
client = httpx.AsyncClient()

查阅 httpx 源码你会发现:

  • max_connections 默认是 100,max_keepalive_connections 默认只有 20。
  • timeout 默认全量覆盖 5.0 秒。

大模型生成长文本往往需要 10 秒甚至数十秒。当并发流达到 100 以上时,后续请求全部在等待空闲连接,5 秒一到立即抛出 httpx.PoolTimeout。同时,如果 LLM 首字延迟(TTFT)稍长,默认的 read timeout 就会把正常请求强行掐死。

优化方案:按流式场景定制 Limits 与 Timeout

import httpx

# 针对大模型 SSE 长连接量身定制
limits = httpx.Limits(
    max_connections=1000,           # 允许的最大并发连接数
    max_keepalive_connections=200, # 保持活跃的连接池容量
    keepalive_expiry=30.0          # 保活过期时间
)

timeout = httpx.Timeout(
    connect=5.0,    # 建立连接必须快
    read=None,      # 关键:大模型流式生成不要限制全局读取超时!write=10.0,     # 发送 Prompt 超时
    pool=10.0       # 排队拿连接的最大容忍时间
)

client = httpx.AsyncClient(limits=limits, timeout=timeout)

第二坑:客户端主动断开,上游 GPU 仍在空转(算力盗刷)

用户在前端点击“停止生成”,或者直接关掉浏览器标签页时,Web 客户端与 FastAPI 网关之间的 HTTP 连接已经切断。但 默认情况下,FastAPI 不会自动终止与上游推理集群(vLLM/Ollama)的通信!

上游模型服务还在傻乎乎地拼命跑 KV Cache 和 Decode 计算,把生成的 Token 塞进操作系统的 Socket 缓冲区,直到整个长文本推理完毕。在公司内部,这意味着有限的 GPU 算力被大量无效请求霸占。

要解决这个问题,网关必须感知下游客户端的存活状态:只要前端断开,立即打断生成器,让上游连接触发关闭回调。

from fastapi import Request

async def stream_generator(upstream_response, request: Request):
    try:
        async for chunk in upstream_response.aiter_raw():
            # 每次准备推流前,先检查客户端是否断连
            if await request.is_disconnected():
                # 记录日志,主动跳出循环
                break
            yield chunk
    finally:
        # 无论正常结束还是客户端断连,务必关停上游流
        await upstream_response.aclose()

第三坑:aiter_lines() 导致的内存悄悄暴涨

在解析上游的 data: {...} 数据格式做鉴权或 Token 计费时,大家通常喜欢用 response.aiter_lines()。如果你压测过就会发现,当上游偶尔吐出包含巨大上下文、无换行符的 base64 图像数据或超长 JSON 时,aiter_lines() 内部会无限扩大 buffer 直到找到换行符为止,极易引发内存峰值抖动。

纯代理网关最安全、最高效的做法是按字节块直接转发(aiter_raw()),不做多余的字符串解码;只有需要做中间件审计(如违禁词过滤、敏感词拦截)时,才结合轻量级增量缓冲区做解析。

生产级代码实现:FastAPI + httpx 流式代理

结合 FastAPI 的 lifespan 生命周期管理连接池,并兼顾流式异常处理的完整实现如下:

import contextlib
from fastapi import FastAPI, Request, HTTPException
from fastapi.responses import StreamingResponse
import httpx

# 1. 生产环境连接池生命周期管理
@contextlib.asynccontextmanager
async def lifespan(app: FastAPI):
    limits = httpx.Limits(max_connections=500, max_keepalive_connections=100)
    timeout = httpx.Timeout(connect=5.0, read=None, write=10.0, pool=5.0)
    app.state.http_client = httpx.AsyncClient(limits=limits, timeout=timeout)
    yield
    await app.state.http_client.aclose()

app = FastAPI(title="LLM-Gateway", lifespan=lifespan)

UPSTREAM_API_URL = "http://192.168.1.100:8000/v1/chat/completions"

@app.post("/v1/chat/completions")
async def chat_proxy(request: Request):
    client: httpx.AsyncClient = request.app.state.http_client
    body = await request.body()
    
    # 过滤可能导致上游混乱的特定 Header
    headers = {k: v for k, v in request.headers.items() 
        if k.lower() not in ("host", "content-length")
    }

    # 构建上游流式请求
    upstream_req = client.build_request(
        method="POST",
        url=UPSTREAM_API_URL,
        content=body,
        headers=headers
    )
    
    try:
        upstream_resp = await client.send(upstream_req, stream=True)
    except httpx.ConnectError:
        raise HTTPException(status_code=502, detail="Upstream LLM server unreachable")
    except httpx.PoolTimeout:
        raise HTTPException(status_code=503, detail="Gateway connection pool exhausted")

    # 2. 具备断连感知的安全推流器
    async def event_publisher():
        try:
            async for chunk in upstream_resp.aiter_raw():
                if await request.is_disconnected():
                    break
                yield chunk
        finally:
            # 关键:确保 upstream 连接在断开后及时回收
            await upstream_resp.aclose()

    return StreamingResponse(event_publisher(),
        status_code=upstream_resp.status_code,
        headers={"Content-Type": upstream_resp.headers.get("content-type", "text/event-stream"),
            "Cache-Control": "no-cache",
            "X-Accel-Buffering": "no"  # 若前面套了 Nginx,务必关闭反代缓冲
        }
    )

实测压测对比:改造前后的数据表现

我们在 4 核 8G 的虚拟机上,使用 Locust 发起 300 并发压测,模拟每秒生成 20 个 Token、持续 30 秒的交互:

  • 未经优化的初始版本:压测至第 45 秒,并发稳定在 120 左右时开始报 PoolTimeout,内存从 90MB 飙到 2.8GB,CPU 占用 100%(大部分耗在无意义的连接建立与断开)。
  • 调优连接池与断开感知后:并发稳定支撑 300 流无报错,单进程内存始终压在 120MB~140MB 之间,无任何线性膨胀现象。当在压测客户端强行中断 50% 的长连接时,上游推理卡的利用率在 200ms 内出现断崖式回落,证明断连感知机制彻底生效。

老司机的避坑经验清单

最后提两点容易忽略的环境问题:

  1. Nginx 缓冲层(Buffering):如果网关前还顶着一层 Nginx,必须在 location 块中设置 proxy_buffering off;,并在响应头带上 X-Accel-Buffering: no,否则 Nginx 会试图把上游数据攒成一个大包再发,导致前端出现“卡顿很久然后文字突然全屏崩出来”的假死现象。
  2. Uvicorn 启动参数:生产部署建议直接使用 uvicorn.workers.UvicornWorker 配合 Gunicorn,单 worker 依靠 uvloop 足以处理数千长连接,不需要开太多 worker 导致连接池资源被无谓切碎。
正文完
 0
评论(没有评论)