实战踩坑:用 FastAPI+SSE 打造低延迟 AI Agent 流式代理与高并发背压治理

11次阅读
没有评论

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

最近在为团队的一套自动化 Agent 平台重构流式网关。很多朋友在用 Python 做大模型流式(Streaming)接口时,以为写个 yield,外面包一层 StreamingResponse(media_type="text/event-stream") 就万事大吉了。然而真正放到稍有并发的场景跑压测,或者遭遇网络较差的客户端频繁中断时,各种诡异问题接踵而至:后端与 vLLM/Ollama 的连接迟迟不释放、长耗时 Tool Call 导致 HTTP 连接挂死、客户端取消请求后 Python 协程依然在空跑甚至内存泄漏。

本文不扯抽象架构,直接记录我从踩坑到压测稳定的全过程,手把手拆解如何用 FastAPI、httpx.AsyncClient 与 SSE 打造一个稳健、防挂死、带背压调优的高并发大模型代理服务。

一、为什么简单的 StreamingResponse 会让生产环境崩掉?

在默认情况下,FastAPI(依托 Starlette)处理异步生成器时,并不会主动监听底层 TCP Socket 的每一丝风吹草动。下面这段代码在技术博客里极为常见,但放到生产就是定时炸弹:

# 常见但有严重隐患的初级写法
@app.post("/v1/chat")
async def chat_stream(req: ChatRequest):
    async def event_generator():
        # 调用下游本地模型(如 vLLM / SGLang / Ollama)async with httpx.AsyncClient() as client:
            async with client.stream("POST", upstream_url, json=req.dict()) as resp:
                async set chunk in resp.aiter_text():
                    yield f"data: {chunk}\n\n"
    return StreamingResponse(event_generator(), media_type="text/event-stream")

这段代码有三个致命缺陷:

  1. 客户端主动断开无感:当用户点击浏览器上的“停止生成”按钮时,浏览器切断 TCP 连接。如果生成器内部没有捕获 http.disconnect 信号,协程会继续向下游拉取 Tokens,白白消耗 GPU 推理算力。
  2. 临时创建连接池的开销灾难:每次请求都 async with httpx.AsyncClient(),完全丧失了底层 TCP 连接复用机制,压测 QPS 刚过 100,系统直接被 TIME_WAIT 拖垮。
  3. Tool Calling 停顿超时:Agent 在流式输出过程中触发本地工具(如执行一段耗时 15 秒的数据查询脚本)时,SSE 通道无任何数据传输,中间层的 Nginx、云厂商 SLB 会直接判定 504 Gateway Timeout 掐断连接。

二、工业级流式网关实现:连接复用与生命周期绑定

解决上述问题的核心是三点:全局单例连接池、基于 request.is_disconnected() 的协作式中断轮询、以及心跳保活机制。

import asyncio
import json
from contextlib import asynccontextmanager
from typing import AsyncGenerator
import httpx
from fastapi import FastAPI, Request, HTTPException
from fastapi.responses import StreamingResponse

# 1. 全局持久化连接池,调优连接数与超时
@asynccontextmanager
async def lifespan(app: FastAPI):
    limits = httpx.Limits(max_keepalive_connections=100, max_connections=500, keepalive_expiry=30.0)
    # 大模型长文本生成,写入和读取 timeout 必须放宽,但 connect timeout 需从严
    timeout = httpx.Timeout(connect=5.0, read=120.0, write=10.0, pool=5.0)
    app.state.client = httpx.AsyncClient(limits=limits, timeout=timeout)
    yield
    await app.state.client.aclose()

app = FastAPI(lifespan=lifespan)

UPSTREAM_ENGINE_URL = "http://127.0.0.1:8000/v1/chat/completions"

async def stream_with_cancellation(
    client: httpx.AsyncClient, 
    payload: dict, 
    request: Request
) -> AsyncGenerator[str, None]:
    headers = {"Content-Type": "application/json"}
    
    try:
        async with client.stream("POST", UPSTREAM_ENGINE_URL, json=payload, headers=headers) as upstream_resp:
            if upstream_resp.status_code != 200:
                err_content = await upstream_resp.aread()
                yield f"event: error\ndata: {json.dumps({'status': upstream_resp.status_code,'detail': err_content.decode()})}\n\n"
                return

            # 使用 aiter_lines() 避免底层 TCP 黏包导致的 SSE 数据帧切片错误
            async for line in upstream_resp.aiter_lines():
                # 检查客户端是否已经关闭连接(如用户关闭页面或主动取消)if await request.is_disconnected():
                    # 极其关键:在此处中断即可立刻触发 upstream_resp 的退出上下文,下游推理请求被 abort
                    break

                if not line:
                    continue

                # 原样透传 SSE 数据帧
                yield f"{line}\n\n"

    except httpx.ReadTimeout:
        yield f"event: error\ndata: {json.dumps({'error':'Upstream LLM engine read timeout'})}\n\n"
    except asyncio.CancelledError:
        # 协程被强制取消时的防御性捕获
        pass

@app.post("/v1/agent/chat")
async def agent_chat_endpoint(request: Request):
    payload = await request.json()
    
    # 强制将请求重置为 stream 模式
    payload["stream"] = True

    client: httpx.AsyncClient = request.app.state.client
    generator = stream_with_cancellation(client, payload, request)
    
    return StreamingResponse(
        generator, 
        media_type="text/event-stream",
        headers={
            "Cache-Control": "no-cache",
            "Connection": "keep-alive",
            "X-Accel-Buffering": "no"  # 禁用 Nginx 缓冲,让数据即刻推送到端
        }
    )

三、防断流关键:针对慢速工具调用的 Ping 帧心跳注入

在 Agent 场景下,大模型经常调用各种复杂的 Python 脚本、矢量数据库或外部 API。当模型决定 tool_calls 并等待本地运行环境返回时,输出流会彻底陷入静默。很多网关和前端反向代理(如 Traefik、Cloudflare、Nginx)默认 proxy_read_timeout 是 60 秒,一旦工具执行耗时稍长,链路直接报废。

解决策略是在后端推流管道中引入 asyncio.wait 并行心跳机制:如果在指定秒数(比如 3 秒)内没有新 Token 生成,就主动发送一条 SSE 注释行(以冒号 : 开头)作为 Keep-Alive 脉冲:

async def stream_with_heartbeat(
    client: httpx.AsyncClient, 
    payload: dict, 
    request: Request,
    heartbeat_interval: float = 3.0
) -> AsyncGenerator[str, None]:
    async with client.stream("POST", UPSTREAM_ENGINE_URL, json=payload) as upstream_resp:
        line_gen = upstream_resp.aiter_lines().__aiter__()
        
        while True:
            if await request.is_disconnected():
                break

            try:
                # 设定心跳超时竞争
                line = await asyncio.wait_for(line_gen.__anext__(), timeout=heartbeat_interval)
                if line:
                    yield f"{line}\n\n"
            except asyncio.TimeoutError:
                # 超过 3 秒未产生新数据,向客户端发送 SSE 注释帧保活
                yield ": keep-alive-pulse\n\n"
            except StopAsyncIteration:
                break

由于标准的 SSE 协议规范中明确规定以 : 开头的行为注释帧,前端 EventSource 会自动丢弃该行而不会污染渲染组件,这就能在零侵入业务前端的情况下让整条连接持续维持存活状态。

四、压测调优:Nginx 与 Linux 内核背压治理

代码层面搞定后,真正的战场在 Linux 系统和反向代理层。我在单台 16 核 32G 机器上用 locust 发起 1000 并发流式测试时,曾遇到服务端内存疯狂吃满的现象。排查发现是 背压(Backpressure)失衡 导致的:由于客户端网络延迟较大(如 4G 移动端),接收数据的速度远远慢于后端推理引擎吐 Token 的速度,导致未消费的包被堆积在 Python 进程的内存队列中。

1. 反向代理层(Nginx)关键参数

必须关闭代理缓冲,同时放开客户端保活限制:

location /v1/agent/ {
    proxy_pass http://127.0.0.1:8080;
    proxy_http_version 1.1;
    proxy_set_header Connection "";
    
    # 彻底关闭缓冲,避免 Token 在 Nginx 内部蓄水池攒满才向前端推送
    proxy_buffering off;
    proxy_cache off;
    
    # 将读写超时放宽至 10 分钟以适应复杂的长推理链路
    proxy_read_timeout 600s;
    proxy_send_timeout 600s;
    chunked_transfer_encoding on;
}

2. 系统级网络参数优化

在 /etc/sysctl.conf 中调优发送缓冲区和文件描述符:

# 调大 TCP 缓冲区范围:min default max
net.ipv4.tcp_wmem = 4096 16384 4194304
net.ipv4.tcp_rmem = 4096 87380 4194304

# 开启 TCP 快速回收与连接重用
net.ipv4.tcp_tw_reuse = 1

# 调高系统最大打开文件句柄
fs.file-max = 2097152

五、实测性能与选型建议

针对 500 并发模拟客户端的实测场景,对比了不同实现模式下的资源占用与稳定性指标(硬件:Ubuntu 22.04 LTS, 16C/32G):

方案实现 500 并发峰值内存 客户端取消响应延迟 弱网环境掉线率
裸写 StreamingResponse (每次新建 Client) 2.8 GB+ (频繁出现 OOM) 无感知(后端持续空转至结束) 18.4% (超时断流)
FastAPI + 全局 Client + 基础 SSE 480 MB 偶发泄漏(依赖 TCP 垃圾回收) 7.2%
本文方案 (心跳注入 + 协作式中断) 210 MB < 50ms (即时感知释放) 0.05%

总结与避坑心得

在大模型工程落地中,流式传输早已不是“玩具代码”阶段。总结下来,这套生产级代理方案核心就在三件事:

  1. 别让后端当冤大头:一定要捕获 request.is_disconnected(),用户点停就立刻掐断上游请求,省下的算力都是白花花的银子;
  2. 连接池必须单例并合理调参:流式请求占用长连接时间远超传统 REST 接口,连接池的 max_connections 要与底层推理服务的负载上限严格对齐;
  3. 用好心跳注释帧:针对 Agent 的长耗时 Tool Call 必须定时喂保活脉冲,把网络断流的隐患掐死在萌芽状态。

建议各位在搭建自己的 AI Agent 路由层时,不要直接依赖臃肿且黑盒的现成框架,先用这套精简的 FastAPI 方案把网络链路调顺,你会发现后期的排障和压测都会顺畅得多。

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