实战踩坑:基于 FastAPI 与 Asyncio 搭建高并发 LLM SSE 流式代理服务

9次阅读
没有评论

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

上个月我给团队内部的聚合 AI 网关做压测,当并发打到 300 以上时,基于传统 requests + 同步多线程包装的 Python 服务内存直接被吃爆,P99 延迟飙升到 8 秒开外,甚至大量请求直接返回 502。究其原因,大模型生成的长连接与常规 REST API 的短连接有着本质区别:一个包含思考链(Thinking Process)的长文本请求动辄耗时 15 到 30 秒,传统并发模型在应对数以千计的 SSE(Server-Sent Events)挂起连接时,极易发生协程雪崩和内存泄漏。

本文抛弃所有花架子架构,分享我用 FastAPI + httpx.AsyncClient 重构整个流式中继服务的实操细节,重点解决连接池耗尽、客户端提前断连导致的显卡算力浪费,以及生产环境中的异步背压控制。

一、长连接流式代理的三个致命暗礁

很多同学在写大模型代理时,代码通常长这样:写一个 StreamingResponse,在生成器里循环 await response.aiter_bytes(),然后 yield 出去。这在本地开发或几个人使用时完全没问题,但一旦放到生产环境,马上会遭遇三个深坑:

  1. 客户端断连感知失效(幽灵推理):用户在 WebUI 上点击“停止生成”或直接关掉浏览器标签页,如果你的代理层没有显式捕获客户端的断开事件(Client Disconnect),后端的 vLLM/Ollama 或商业 API 仍会傻傻地把几千个 Token 生成完毕。这不仅占用宝贵的 GPU 显存,还会白白消耗上游 API 的计费额度。
  2. 全局异步 Client 连接池枯竭:随意在每个请求中实例化 httpx.AsyncClient() 会导致严重的 TCP 连接建立开销与 TIME_WAIT 积压;而单例 Client 如果没有合理配置连接池上限(max_connections 与 max_keepalive_connections),瞬间高并发就会直接卡死在获取连接的信号量上。
  3. 慢客户端导致的内存膨胀(缺乏背压机制):上游大模型生成速度如果达到 80 tokens/s,而下游客户端处于高延迟弱网环境(例如移动端),代理层如果无脑在内存中缓存缓冲区数据,并发上来后服务内存就会呈线性暴涨。

二、生产级架构与核心代码落地

为了解决上述问题,我们需要将 FastAPI 升级为带生命周期管理、连接感知与背压保护的异步中继服务。以下是我在生产环境中经过压测验证的核心代理模块代码:

import asyncio
import logging
from contextlib import asynccontextmanager
from typing import AsyncGenerator

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

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("llm_proxy")

# 1. 严格配置全局连接池参数
UPSTREAM_TIMEOUT = httpx.Timeout(
    connect=5.0,
    read=60.0,
    write=5.0,
    pool=10.0
)
LIMITS = httpx.Limits(
    max_connections=1000,
    max_keepalive_connections=200,
    keepalive_expiry=30.0
)

# 2. 生命周期管理:单例复用 AsyncClient
class GlobalState:
    client: httpx.AsyncClient = None

state = GlobalState()

@asynccontextmanager
async def lifespan(app: FastAPI):
    # 服务启动时初始化
    state.client = httpx.AsyncClient(
        limits=LIMITS,
        timeout=UPSTREAM_TIMEOUT,
        http2=True  # 启用 HTTP/2 提升多路复用效率
    )
    logger.info("Global httpx.AsyncClient initialized.")
    yield
    # 优雅停机
    await state.client.aclose()
    logger.info("Global httpx.AsyncClient closed.")

app = FastAPI(lifespan=lifespan)

# 3. 核心:带断开检测与背压控制的流式转发器
async def stream_generator(
    request: Request,
    upstream_url: str,
    payload: dict,
    headers: dict
) -> AsyncGenerator[bytes, None]:
    req = state.client.build_request("POST", upstream_url, json=payload, headers=headers)
    
    try:
        upstream_resp = await state.client.send(req, stream=True)
        upstream_resp.raise_for_status()
    except httpx.HTTPStatusError as exc:
        logger.error(f"Upstream HTTP error: {exc.response.status_code}")
        yield f"data: {{\"error\": \"Upstream error {exc.response.status_code}\"}}\n\n".encode("utf-8")
        return
    except Exception as exc:
        logger.error(f"Connection failed: {str(exc)}")
        yield b"data: {\"error\": \"Upstream connection failed\"}\n\n"
        return

    # 声明消费与健康状态
    try:
        async for chunk in upstream_resp.aiter_bytes():
            # 关键:检查客户端连接是否已经断开
            if await request.is_disconnected():
                logger.warning("Downstream client disconnected. Terminating upstream stream.")
                break
            
            # 正常回传 SSE Chunk
            yield chunk
            # 极低开销让出事件循环,避免单一长流霸占协程
            await asyncio.sleep(0)
            
    except asyncio.CancelledError:
        logger.info("Stream task was explicitly cancelled.")
        raise
    finally:
        # 无论下游异常退出还是正常结束,强制释放上游流连接
        await upstream_resp.aclose()
        logger.debug("Upstream response stream safely closed.")

@app.post("/v1/chat/completions")
async def chat_proxy(request: Request):
    try:
        body = await request.json()
    except Exception:
        raise HTTPException(status_code=400, detail="Invalid JSON body")

    upstream_url = "http://127.0.0.1:8000/v1/chat/completions" # 替换为真实模型后端地址
    headers = {"Authorization": request.headers.get("Authorization", ""),"Content-Type":"application/json"
    }

    # 针对非流式请求走常规快路径
    if not body.get("stream", False):
        resp = await state.client.post(upstream_url, json=body, headers=headers)
        return resp.json()

    # 流式请求转发
    return StreamingResponse(stream_generator(request, upstream_url, body, headers),
        media_type="text/event-stream",
        headers={
            "Cache-Control": "no-cache",
            "Connection": "keep-alive",
            "X-Accel-Buffering": "no"  # 禁用 Nginx 默认缓冲,防止流被积压
        }
    )

三、生产部署必调的底层细节

1. 反向代理层(Nginx)的缓冲黑洞

很多小伙伴在本地测试流式打字机效果很流畅,一上生产通过 Nginx 之后,文字变成几百字一坨一坨地吐出来。这是因为 Nginx 默认开启了 proxy_buffering。除了在响应头中加入 X-Accel-Buffering: no 之外,Nginx 的核心配置务必明确关闭缓冲:

location /v1/chat/completions {
    proxy_pass http://127.0.0.1:8080;
    proxy_set_header Connection '';
    proxy_http_version 1.1;
    proxy_buffering off;
    proxy_cache off;
    chunked_transfer_encoding on;
    # SSE 长连接超时时间建议拉长至 10~15 分钟
    proxy_read_timeout 600s;
}

2. Uvicorn 进程与事件循环选型

在 Linux 生产环境中,尽量使用 uvloop 替代 Python 默认的 asyncio 事件循环,单核 IO 吞吐量能提升 25%~40%。在 Docker 容器或 systemd 中启动命令建议配置为:

uvicorn main:app --host 0.0.0.0 --port 8080 --workers 4 --loop uvloop --http httptools --limit-concurrency 2000

注意 --limit-concurrency 2000 这个参数,它会在系统过载时主动返回 503 保护服务本身,防止系统直接 OOM 崩溃。

四、实测压测数据对比

我在一台 8 核 16G 的测试机上,使用 Locust 模拟 500 个并发长连接用户持续请求 15 秒输出的大模型流式响应,对比了优化前后各项指标:

架构方案 内存占用(稳态) 首字延迟 (TTFT P99) 客户端主动断开后上游浪费率 吞吐稳定性
旧版:同步多线程 + Requests 3.2 GB(频繁上升) 4800 ms 100%(一直空跑到结束) 300 并发开始大量 502
新版:Asyncio + 连接池 + 断连感知 420 MB 380 ms < 2%(断开后 50ms 内熔断) 1200 并发无报错平稳运行

实测结果非常直观:通过合理的单例连接池复用与 request.is_disconnected() 主动探测,服务常驻内存降低了近 87%,且消除了后端模型节点的无意义推理消耗。

总结与避坑心得

写大模型应用与传统 Web 业务最大的不同在于 长时占用(Duration)取代了短频快(Throughput)。你在写代码时时刻要问自己三个问题:

  • 如果客户端突然掉线了,我的上游协程有没有被及时 kill 掉?
  • 如果下游网络极慢,我的内存是不是在无限缓冲数据?
  • 每一次 HTTP 请求发起时,底层的 TCP 连接是不是从统一管控的连接池中借调的?

把这三个问题处理干净,即使是一台入门级云服务器,也能稳稳支撑起数千路并发的大模型 SSE 流式网关。

===END===

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