Skip to content

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 协议,是一种单向推送机制——服务器向客户端持续发送数据,但客户端不能通过同一连接反向发送数据。它的工作方式很简单:

  1. 客户端通过 HTTP 请求建立连接
  2. 服务器保持连接打开,持续以 data: <内容>\n\n 格式发送事件
  3. 连接断开时,浏览器会自动重连

SSE 的数据格式遵循 W3C 规范,每条消息以 data: 开头,以 \n\n(两个换行)结尾:

data: {"content": "你"}\n\n
data: {"content": "好"}\n\n
data: [DONE]\n\n

浏览器原生提供 EventSource API 来消费 SSE,也可以用 fetch + ReadableStream 手动解析。

WebSocket

WebSocket 是一种双向通信协议,客户端和服务器可以在同一个连接上互发消息。它的工作流程是:

  1. 客户端发起 HTTP 升级请求(Upgrade: websocket
  2. 握手成功后,连接从 HTTP 切换为 WebSocket 协议(ws://
  3. 双方可以在任意时刻互相发送消息
  4. 需要手动处理重连逻辑

两者的对比

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 流式输出接口,代码逐行注释:

python
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 发过来,你必须在流的各处拼接它。

python
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 原生实现:

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 实现(带流式状态):

jsx
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 等推理模型支持将"思考过程"也流式输出。这让用户能看到模型是如何推理的,而不仅仅是最终答案。

python
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 字段决定把内容渲染到思考区域还是回答区域:

javascript
if (parsed.type === 'reasoning') {
    setReasoning(prev => prev + parsed.content);  // 更新思考区
} else if (parsed.type === 'answer') {
    setAnswer(prev => prev + parsed.content);     // 更新回答区
}

中断与取消

用户可能想中途停止生成(比如发现回答方向不对)。这需要在服务端实现取消机制:

python
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 完整服务示例

将前面的内容整合成一个可运行的完整服务,包含取消功能:

python
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:

python
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_content

9.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 WebSocketSSE 单向推送适合 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 的准确性与效率,以及如何基于评估结果进行迭代优化。从"能跑"到"跑得好",评估是不可或缺的桥梁。