实战踩坑:Python构建高并发LLM流式代理网关与背压调优

7次阅读
没有评论

共计 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 重构网关服务,记住这三条铁律:

  1. 永远不要在每次请求内部随手 httpx.AsyncClient(),必须在应用生命周期全局维护连接池。
  2. 必须在流生成器内部主动轮询 request.is_disconnected(),把取消信号传递给上游,这是保住 GPU 显存的核心手段。
  3. 全链路关闭数据缓冲,包括网关自身的响应配置和前端各类 Nginx/Traefik 负载均衡器。

用这套骨架搭建的网关,代码量少且高度可控,不用引入 Java/Go 那些重型网关框架,小团队单台服务器即可轻松扛住数千路实时的大模型长连接交互。

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