在 Agent 流式输出场景中,如果用户中途关闭了浏览器,服务端应该如何感知并停止 LLM 调用,避免产生不必要的 Token 费用?
这个问题在我们做智能客服 Agent 的时候特别关键——用户聊着聊着突然关掉网页,后端要是不管不顾继续生成 token,不仅浪费钱,还可能把对话上下文搞脏。下面我按 感知断开 → 取消 LLM 调用 → 实际落地 这条线来说。
🔸 第一步:怎么知道用户已经离开了?
不管你是用 SSE(Server-Sent Events)还是 WebSocket,底层都是 HTTP/1.1 的长连接或 HTTP/2 的流。客户端关闭页面时,浏览器会主动关闭 TCP 连接,服务端可以从几个层面感知到:
1️⃣ 检测写入失败(最直接)
在迭代生成器往响应体写数据的时候,如果连接断开,下一次 write 或 send 会触发 ConnectionResetError 或 BrokenPipeError。
try:
for chunk in stream:
yield chunk
except GeneratorExit:
# FastAPI 的 StreamingResponse 迭代器被取消时会触发
cleanup()
2️⃣ 监听底层请求对象的状态(Starlette/FastAPI)
async def event_generator(request: Request):
for chunk in stream:
if await request.is_disconnected():
# 立即停止 LLM 调用
break
yield f"data: {chunk}\n\n"
FastAPI 的 Request.is_disconnected() 会检查底层 TCP 连接是否还活着,它其实是轮询底层 ASGI 的 transport.is_closing()。
3️⃣ 利用 ASGI 的取消机制
在 ASGI 应用里,当客户端断开时,Uvicorn 会往协程里抛出一个 asyncio.CancelledError。你可以在生成器里捕获它:
async def safe_stream(request: Request):
try:
# 这里调用 LLM
async for token in llm.generate():
yield f"data: {token}\n\n"
except asyncio.CancelledError:
# 客户端已断开
await llm.cancel() # 中断 LLM 请求
raise
🔹 第二步:如何停止正在进行的 LLM 调用?
LLM 的 API 一般都是基于 HTTP 的请求,部分 SDK 支持流式读取。想要省 token,就得在连接断开时主动取消请求:
① 使用 OpenAI SDK 的取消机制
from openai import AsyncOpenAI
client = AsyncOpenAI()
async for token in await client.chat.completions.create(
model="gpt-4",
messages=[...],
stream=True
):
yield token.choices[0].delta.content
如果迭代过程中发生 CancelledError,你需要关闭底层的 HTTP 连接,否则 SDK 可能还在缓冲数据。我们可以这样:
try:
stream = await client.chat.completions.create(..., stream=True)
async for token in stream:
yield token.choices[0].delta.content
finally:
await stream.response.aclose() # 确保释放连接
aclose()会优雅地关闭 HTTP 流,阻止继续接收数据,也就不会消耗更多 token(API 端也会停止生成)。
② 调用 LLM 时加上 request 状态检查
把健康检查包装成一个工具,在每次生成 token 前都查一下:
async def stream_with_check(request: Request, llm_stream):
async for token in llm_stream:
if await request.is_disconnected():
# 强行关闭正在迭代的流
await llm_stream.aclose()
break
yield token
⚠️ 但这招有延迟,可能已经多收了几十个 token。最敏捷的做法还是依赖 CancelledError + finally 清理。
🛠️ 第三步:落地架构(我们用过的模式)
我们给 Agent 流式接口做了一层统一的生命周期管理器:
Client 请求 ──► FastAPI路由
│
├─ 创建 LLM 请求任务 (Task)
│
├─ 在 StreamingResponse 的生成器里:
│ • 检测 is_disconnected 或 捕获 CancelledError
│ • 取消 LLM Task
│ • 清理上下文、锁、日志
│
└─ 返回 SSE 流
代码骨架:
from fastapi import FastAPI, Request
from fastapi.responses import StreamingResponse
import asyncio
app = FastAPI()
async def agent_stream(prompt: str, request: Request):
llm_task = None
try:
# 启动 LLM 协程任务
llm_task = asyncio.create_task(run_llm(prompt))
while not llm_task.done():
if await request.is_disconnected():
llm_task.cancel() # 取消任务
break
# 从队列获取 token 并 yield
token = await queue.get()
yield f"data: {token}\n\n"
queue.task_done()
except asyncio.CancelledError:
if llm_task:
llm_task.cancel()
raise
finally:
# 确保 LLM 连接关闭
if llm_task and not llm_task.done():
llm_task.cancel()
我们在 run_llm 里封装了 OpenAI 调用和异常处理,任务被取消后会进入 except asyncio.CancelledError 清理 HTTP 连接。