9.6 流式输出 Streaming
承前:从成本控制到用户体验
在上一节中,我们讨论了成本控制——如何通过缓存、模型路由、Token 预算等手段让 Agent 的运行开销降到最低。然而,当一个问题需要 LLM 生成数百甚至上千个 Token 的回答时,用户体验往往不取决于成本,而取决于等待感。想象一下:用户发了一个问题,然后盯着空白屏幕等了 15 秒才看到答案——即使答案再好,体验也已经打了折扣。
流式输出(Streaming)正是解决这个问题的核心手段。它让用户在模型生成的同时就能看到逐字、逐句的输出,把"等待"变成"阅读",把焦虑变成惊喜。
9.6.1 什么是流式输出:水龙头放水 vs 等满桶
要理解流式输出,不妨用一个生活类比。
等满桶(非流式输出): 你想接一桶水。你打开水龙头,然后把水桶放到下面等——等 10 分钟,桶装满了,你才把桶提走。在这 10 分钟里,你什么都看不到,只能干等。这就是传统的非流式 API 调用:客户端发送请求,服务器生成完整回复,最后一次性返回。如果生成花了 15 秒,用户就要盯着空白界面等 15 秒。
水龙头放水(流式输出): 同样接一桶水,但这次你打开水龙头后立刻就能看到水在流入——第一滴水入桶,第二滴紧跟其后,水桶逐渐变满。你不需要等桶装满才能看到进展,每一滴都在即时反馈。这就是流式输出:服务器一边生成、一边把每个 Token(或 chunk)推送给客户端,用户在生成的过程中就能实时看到文字。
| 维度 | 等满桶(非流式) | 水龙头放水(流式) |
|---|---|---|
| 返回方式 | 全部生成完才返回 | 边生成边返回 |
| 首字延迟 | = 总生成时间 | ≈ 第一个 Token 的时间 |
| 用户感知 | 长时间空白等待 | 即时反馈、逐字显示 |
| 典型场景 | 后台批处理 | 对话式 AI、实时写作 |
| 实现复杂度 | 简单 | 需要流式协议支持 |
关键区别在于首字延迟(Time to First Token, TTFT)。 在非流式模式下,TTFT 等于整个生成时间;在流式模式下,TTFT 通常只有几百毫秒——模型"想好"第一个词就能送出来。
9.6.2 SSE vs WebSocket:两种流式通道
要实现"水龙头放水"的效果,服务器需要一个能持续向客户端推送数据的通道。目前主要有两种选择:SSE 和 WebSocket。
SSE(Server-Sent Events,服务器推送事件)
SSE 基于 HTTP 协议,是一种单向推送机制——服务器向客户端持续发送数据,但客户端不能通过同一连接反向发送数据。它的工作方式很简单:
- 客户端通过 HTTP 请求建立连接
- 服务器保持连接打开,持续以
data: <内容>\n\n格式发送事件 - 连接断开时,浏览器会自动重连
SSE 的数据格式遵循 W3C 规范,每条消息以 data: 开头,以 \n\n(两个换行)结尾:
data: {"content": "你"}\n\n
data: {"content": "好"}\n\n
data: [DONE]\n\n浏览器原生提供 EventSource API 来消费 SSE,也可以用 fetch + ReadableStream 手动解析。
WebSocket
WebSocket 是一种双向通信协议,客户端和服务器可以在同一个连接上互发消息。它的工作流程是:
- 客户端发起 HTTP 升级请求(
Upgrade: websocket) - 握手成功后,连接从 HTTP 切换为 WebSocket 协议(
ws://) - 双方可以在任意时刻互相发送消息
- 需要手动处理重连逻辑
两者的对比
SSE (Server-Sent Events) WebSocket
───────────────────────── ──────────────────
单向:服务器 → 客户端 双向:服务器 ↔ 客户端
基于 HTTP(普通请求) 独立协议(ws:// / wss://)
浏览器原生 EventSource 需要 WebSocket API
断开后自动重连 需手动实现重连
穿过 HTTP 代理/CDN 无障碍 代理可能不支持或需特殊配置
消息格式简单(data: ...) 消息格式自由(需自定义协议)
连接数限制(浏览器每域名约 6 个) 无特殊限制
适合:LLM 文本流式输出、通知推送 适合:实时聊天、协同编辑、游戏为什么 LLM 流式输出通常选择 SSE? 因为 LLM 的输出本质是单向的——模型生成文本推给用户,用户不需要在同一个数据通道里实时发消息(发消息可以另外发一次普通 HTTP 请求)。SSE 更简单、兼容性更好,能直接穿过 Nginx、CDN 等常见的 HTTP 基础设施,不需要额外的协议升级。只有在需要真正的双向实时通信时(比如协作编辑文档),才值得用 WebSocket。
9.6.3 FastAPI 实现 SSE 流式输出
下面用 FastAPI 实现一个完整的 SSE 流式输出接口,代码逐行注释:
from fastapi import FastAPI # 导入 FastAPI 框架主类
from fastapi.responses import StreamingResponse # 导入流式响应类,用于 SSE 推送
from openai import AsyncOpenAI # 导入 OpenAI 异步客户端
import asyncio # 异步编程标准库
import json # JSON 序列化/反序列化
app = FastAPI() # 创建 FastAPI 应用实例
client = AsyncOpenAI() # 创建 OpenAI 异步客户端,自动读取环境变量中的 API Key
async def generate_stream(messages: list):
"""异步生成器:将 LLM 的流式输出转为 SSE 格式"""
# 调用 OpenAI 接口,设置 stream=True 开启流式模式
stream = await client.chat.completions.create(
model="gpt-4o-mini", # 使用的模型名
messages=messages, # 对话历史
stream=True # 关键参数:开启流式返回
)
# 异步遍历流式响应,每次迭代得到一个 chunk(包含一小段增量内容)
async for chunk in stream:
delta = chunk.choices[0].delta # delta 是本次 chunk 的增量内容
if delta.content: # 如果有文本内容
data = json.dumps({ # 将内容序列化为 JSON 字符串
"content": delta.content, # 本次增量文本
"finish_reason": chunk.choices[0].finish_reason # 结束原因(最后一段才有值)
}, ensure_ascii=False) # 保留中文字符,不转义为 \uXXXX
yield f"data: {data}\n\n" # 按 SSE 格式输出:data: <json> + 两个换行
yield "data: [DONE]\n\n" # 流结束后发送 [DONE] 标记,通知客户端结束
@app.post("/chat/stream") # 定义 POST 路由
async def chat_stream(request: dict): # request 自动解析请求体为 dict
messages = request["messages"] # 提取对话消息列表
return StreamingResponse( # 返回 StreamingResponse,FastAPI 会持续推送生成器产出的数据
generate_stream(messages), # 传入异步生成器作为数据源
media_type="text/event-stream", # 设置 Content-Type 为 SSE 类型
headers={
"Cache-Control": "no-cache", # 禁用缓存,确保每条数据即时推送
"Connection": "keep-alive", # 保持连接不断开
"X-Accel-Buffering": "no", # 关键:禁用 Nginx 的响应缓冲,否则数据会被攒在一起再发
}
)注意
X-Accel-Buffering: no:这是 Nginx 特有的响应头。默认情况下 Nginx 会缓冲响应内容,导致 SSE 数据不能即时推送。设置此头后 Nginx 会逐条转发,这是生产环境中常见的"为什么我的 SSE 不流式"问题的解法。
9.6.4 流式状态管理:分块累积
流式输出的一个挑战是:数据是分块到达的,但你往往需要完整的信息才能做后续处理。比如工具调用(Tool Call)的参数 JSON 可能被拆成好几个 chunk 发过来,你必须在流的各处拼接它。
class StreamManager:
"""管理流式输出的状态,累积分块数据"""
def __init__(self):
self.buffer = "" # 文本缓冲区,累积所有文本内容
self.tool_calls = [] # 工具调用列表,支持流式拼接 JSON 参数
self.current_tool = None # 当前正在构建的工具调用索引
def process_chunk(self, chunk):
"""处理单个 chunk,返回事件描述"""
delta = chunk.choices[0].delta # 取出增量部分
# 情况一:普通文本内容
if delta.content:
self.buffer += delta.content # 把新文本追加到缓冲区
return {"type": "text", "content": delta.content} # 返回增量给前端
# 情况二:工具调用(参数被拆分到多个 chunk 中)
if delta.tool_calls:
for tc in delta.tool_calls:
# 如果是新工具调用(index 超出已有列表长度),先创建占位
if tc.index >= len(self.tool_calls):
self.tool_calls.append({
"id": tc.id or "", # 工具调用 ID(首 chunk 才有)
"function": {"name": "", "arguments": ""} # 函数名和参数(后续逐步拼接)
})
# 补充工具调用 ID(通常只在第一个 chunk 出现)
if tc.id:
self.tool_calls[tc.index]["id"] = tc.id
if tc.function:
if tc.function.name:
self.tool_calls[tc.index]["function"]["name"] += tc.function.name # 追加函数名
if tc.function.arguments:
self.tool_calls[tc.index]["function"]["arguments"] += tc.function.arguments # 追加参数片段
return {"type": "tool_call_building", "data": self.tool_calls[tc.index]}
return {"type": "other"} # 其他类型的 delta(如 role 等),暂时忽略
def finalize(self):
"""流结束后返回完整结果"""
return {
"content": self.buffer, # 完整文本
"tool_calls": self.tool_calls if self.tool_calls else None # 完整工具调用(如有)
}核心思路:process_chunk 在每个 chunk 到达时被调用,负责累积;finalize 在流结束后返回拼好的完整数据。这样前端可以实时显示部分内容,而 Agent 逻辑可以在流结束后拿到完整的工具调用信息继续执行。
9.6.5 前端消费 SSE
流式输出的另一半是前端渲染。服务器在"放水",前端需要"接水"并逐字显示。
JavaScript 原生实现:
async function streamChat(messages) {
// 向后端 SSE 接口发送 POST 请求
const response = await fetch('/chat/stream', {
method: 'POST', // 使用 POST 以便发送消息体
headers: { 'Content-Type': 'application/json' }, // 声明 JSON 格式
body: JSON.stringify({ messages }) // 序列化消息数组
});
const reader = response.body.getReader(); // 获取 ReadableStream 的读取器
const decoder = new TextDecoder(); // 创建 UTF-8 文本解码器
let fullContent = ''; // 累积完整内容
while (true) {
const { done, value } = await reader.read(); // 读取一块数据
if (done) break; // 流结束则退出循环
const text = decoder.decode(value); // 将字节解码为文本
const lines = text.split('\n'); // 按行分割
for (const line of lines) {
if (line.startsWith('data: ')) { // SSE 数据行
const data = line.slice(6); // 去掉 "data: " 前缀
if (data === '[DONE]') return fullContent; // 结束标记
const parsed = JSON.parse(data); // 解析 JSON
fullContent += parsed.content; // 追加增量内容
updateChatUI(fullContent); // 更新界面(逐字显示)
}
}
}
return fullContent;
}React 实现(带流式状态):
function ChatStream() {
const [content, setContent] = useState(''); // 已接收的完整文本
const [isStreaming, setIsStreaming] = useState(false); // 是否正在接收流
async function handleSend(message) {
setIsStreaming(true); // 开始流式接收
setContent(''); // 清空上一次的内容
const response = await fetch('/chat/stream', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ messages: [{ role: 'user', content: message }] })
});
const reader = response.body.getReader(); // 获取流读取器
const decoder = new TextDecoder();
while (true) {
const { done, value } = await reader.read();
if (done) break;
const lines = decoder.decode(value).split('\n');
for (const line of lines) {
if (line.startsWith('data: ') && line.slice(6) !== '[DONE]') {
const { content: chunk } = JSON.parse(line.slice(6));
setContent(prev => prev + chunk); // 使用函数式更新,避免闭包陷阱
}
}
}
setIsStreaming(false); // 流结束
}
return (
<div>
<div className="message">{content}</div>
{isStreaming && <span className="cursor">▊</span>} {/* 流式时显示光标 */}
</div>
);
}React 注意事项:使用
setContent(prev => prev + chunk)而非setContent(content + chunk),因为闭包中的content可能不是最新值。这是流式渲染中最常见的 bug 来源。
9.6.6 高级流式功能
思考过程展示(Reasoning Streaming)
DeepSeek-R1、OpenAI o1 等推理模型支持将"思考过程"也流式输出。这让用户能看到模型是如何推理的,而不仅仅是最终答案。
async def generate_with_reasoning(messages):
"""分离思考过程和最终回答,分别流式推送"""
stream = await client.chat.completions.create(
model="deepseek-reasoner", # 支持推理流式输出的模型
messages=messages,
stream=True
)
async for chunk in stream:
delta = chunk.choices[0].delta
# 推理内容(思考过程)
if hasattr(delta, 'reasoning_content') and delta.reasoning_content:
yield f"data: {json.dumps({'type': 'reasoning', 'content': delta.reasoning_content}, ensure_ascii=False)}\n\n"
# 最终回答内容
if delta.content:
yield f"data: {json.dumps({'type': 'answer', 'content': delta.content}, ensure_ascii=False)}\n\n"
yield "data: [DONE]\n\n"前端可以根据 type 字段决定把内容渲染到思考区域还是回答区域:
if (parsed.type === 'reasoning') {
setReasoning(prev => prev + parsed.content); // 更新思考区
} else if (parsed.type === 'answer') {
setAnswer(prev => prev + parsed.content); // 更新回答区
}中断与取消
用户可能想中途停止生成(比如发现回答方向不对)。这需要在服务端实现取消机制:
import asyncio
class CancellableStream:
"""支持中途取消的流式输出"""
def __init__(self):
self._cancel_event = asyncio.Event() # 创建事件标志,初始为 False
def cancel(self):
"""外部调用此方法触发取消"""
self._cancel_event.set() # 设置事件标志为 True
async def stream_with_cancel(self, messages):
"""带取消检查的流式生成器"""
stream = await client.chat.completions.create(
model="gpt-4o-mini",
messages=messages,
stream=True
)
async for chunk in stream:
if self._cancel_event.is_set(): # 每个 chunk 前检查取消标志
yield "data: [CANCELLED]\n\n" # 发送取消信号
await stream.close() # 关闭底层流,释放资源
break
if chunk.choices[0].delta.content:
data = json.dumps({"content": chunk.choices[0].delta.content}, ensure_ascii=False)
yield f"data: {data}\n\n"9.6.7 完整服务示例
将前面的内容整合成一个可运行的完整服务,包含取消功能:
from fastapi import FastAPI, Request
from fastapi.responses import StreamingResponse
from openai import AsyncOpenAI
import json, asyncio
app = FastAPI()
client = AsyncOpenAI()
active_streams = {} # 全局字典:stream_id -> cancel_event
@app.post("/chat/stream")
async def chat_stream(request: Request):
body = await request.json()
messages = body["messages"]
stream_id = body.get("stream_id", str(id(messages))) # 生成或使用传入的 stream_id
async def event_generator():
cancel_event = asyncio.Event()
active_streams[stream_id] = cancel_event # 注册到全局表,供取消接口调用
try:
stream = await client.chat.completions.create(
model="gpt-4o-mini",
messages=messages,
stream=True
)
async for chunk in stream:
if cancel_event.is_set(): # 检查取消
yield f"data: {json.dumps({'type': 'cancelled'})}\n\n"
break
delta = chunk.choices[0].delta
if delta.content:
yield f"data: {json.dumps({'type': 'text', 'content': delta.content}, ensure_ascii=False)}\n\n"
yield "data: [DONE]\n\n"
finally:
active_streams.pop(stream_id, None) # 清理全局表
return StreamingResponse(
event_generator(),
media_type="text/event-stream"
)
@app.post("/chat/stream/{stream_id}/cancel")
async def cancel_stream(stream_id: str):
"""取消指定流"""
if stream_id in active_streams:
active_streams[stream_id].set() # 触发取消事件
return {"status": "cancelled"}
return {"status": "not_found"}Python 客户端消费 SSE:
import httpx
import json
async def consume_sse(url: str, messages: list):
"""消费 SSE 流式响应并实时打印"""
full_content = ""
async with httpx.AsyncClient(timeout=60.0) as client: # 设置超时避免长时间挂起
async with client.stream(
"POST", url,
json={"messages": messages}
) as response:
async for line in response.aiter_lines(): # 逐行读取流式响应
if line.startswith("data: "):
data = line[6:] # 去掉 "data: " 前缀
if data == "[DONE]":
break # 收到结束标记
chunk = json.loads(data)
if chunk.get("type") == "text":
full_content += chunk["content"]
print(chunk["content"], end="", flush=True) # 实时打印,不换行不缓冲
print() # 最后换行
return full_content9.6.8 常见误区
误区一:SSE 和 WebSocket 可以随意替换
不少开发者认为"既然都是流式,用哪个都行"。实际上两者的适用场景截然不同。SSE 是单向推送,适合 LLM 文本输出这种"只发不收"的场景;WebSocket 是双向通信,适合聊天室、协同编辑。用 WebSocket 做 LLM 流式输出是杀鸡用牛刀——多了协议升级的复杂度,还要自己处理重连。反过来,用 SSE 做实时协作则会发现它无法满足客户端实时回传的需求。
误区二:流式输出能减少总生成时间
流式输出改变的是首字延迟(用户何时看到第一个字),而不是总生成时间(模型生成完所有内容需要多久)。模型该花 10 秒还是 10 秒,流式只是把这 10 秒拆成了"0.3 秒看到第一个字 + 9.7 秒陆续输出",而不是"10 秒后一次性输出全部"。对于需要完整结果才能进行后续处理的场景(比如工具调用),流式不会带来计算上的加速。
误区三:忘记设置 X-Accel-Buffering: no
这是生产环境中最常见的"SSE 不流式"问题。Nginx 默认会缓冲响应,等到一定量才一次性转发。加上 X-Accel-Buffering: no 响应头后,Nginx 会立即转发每个 chunk。如果用了其他反向代理(如 Cloudflare),也要检查是否有类似的缓冲设置。
误区四:前端 React 闭包陷阱
在 React 中直接写 setContent(content + chunk) 时,content 引用的是闭包中的旧值(可能是空字符串),导致后续的追加全部丢失。正确做法是使用函数式更新:setContent(prev => prev + chunk),确保每次更新都基于最新状态。
误区五:取消后不清理资源
用户点击"停止生成"后,如果服务端没有正确关闭底层 LLM 流(await stream.close()),连接会继续消耗资源,甚至继续产生 Token 计费。务必在取消后关闭流并在 finally 块中清理全局状态。
9.6.9 本节小结
| 要点 | 说明 |
|---|---|
| 流式输出的本质 | 边生成边推送,降低首字延迟,提升用户体验 |
| SSE vs WebSocket | SSE 单向推送适合 LLM 输出;WebSocket 双向通信适合协作场景 |
| SSE 数据格式 | data: <json>\n\n,以 [DONE] 结束 |
| FastAPI 实现 | StreamingResponse + 异步生成器 yield |
| 关键响应头 | X-Accel-Buffering: no 禁用 Nginx 缓冲 |
| 状态管理 | StreamManager 累积分块文本和工具调用参数 |
| 推理流式 | DeepSeek-R1 等模型支持 reasoning_content 分离思考与回答 |
| 取消机制 | asyncio.Event + 全局注册表,取消后必须关闭流和清理资源 |
| 前端渲染 | fetch + ReadableStream + TextDecoder,React 中注意闭包陷阱 |
核心一句话: 流式输出把"等待"变成"阅读",是提升 AI 应用体验感的关键技术——就像水龙头放水让你看到水在流,而不是盯着空桶干等。
启后:从系统设计到评估优化
本章我们从系统设计的角度,依次讨论了架构模式、记忆系统、工具编排、多 Agent 协作、成本控制以及流式输出。至此,一个完整的 AI Agent 系统的骨架已经搭好了。但系统搭好只是第一步——它好用吗?它靠谱吗?它够快吗? 这些问题需要通过系统化的评估来回答。
下一章我们将进入第 10 章:评估与优化,讨论如何建立评估体系、设计测试用例、衡量 Agent 的准确性与效率,以及如何基于评估结果进行迭代优化。从"能跑"到"跑得好",评估是不可或缺的桥梁。