共计 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")
这段代码有三个致命缺陷:
- 客户端主动断开无感:当用户点击浏览器上的“停止生成”按钮时,浏览器切断 TCP 连接。如果生成器内部没有捕获
http.disconnect信号,协程会继续向下游拉取 Tokens,白白消耗 GPU 推理算力。 - 临时创建连接池的开销灾难:每次请求都
async with httpx.AsyncClient(),完全丧失了底层 TCP 连接复用机制,压测 QPS 刚过 100,系统直接被TIME_WAIT拖垮。 - 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% |
总结与避坑心得
在大模型工程落地中,流式传输早已不是“玩具代码”阶段。总结下来,这套生产级代理方案核心就在三件事:
- 别让后端当冤大头:一定要捕获
request.is_disconnected(),用户点停就立刻掐断上游请求,省下的算力都是白花花的银子; - 连接池必须单例并合理调参:流式请求占用长连接时间远超传统 REST 接口,连接池的
max_connections要与底层推理服务的负载上限严格对齐; - 用好心跳注释帧:针对 Agent 的长耗时 Tool Call 必须定时喂保活脉冲,把网络断流的隐患掐死在萌芽状态。
建议各位在搭建自己的 AI Agent 路由层时,不要直接依赖臃肿且黑盒的现成框架,先用这套精简的 FastAPI 方案把网络链路调顺,你会发现后期的排障和压测都会顺畅得多。