人类对“等待”的耐心是有限的:当用户发出一个问题,超过 3 秒没有反馈,流失率就会显著上升。LLM 应用因此把“生成中”的状态转化为体验的一部分——流式输出让模型一个字一个字地“打字”,用户感受到的不是等待,而是思考的进行。但流式不是简单地把 HTTP 响应切成小块。它牵扯协议选型、增量解析、中断语义、首 token 延迟(TTFT)、超时重试与可观测性的一整套工程体系。本指南系统覆盖流式架构全景、SSE 与 WebSocket 协议、流式渲染与半流式、超时中断治理以及流式监控,给出可直接落地的完整实践。
一、流式交互的架构全景
1.1 为什么需要流式
LLM 生成的完整响应动辄数百甚至上千 token,串行等待会让体验变得不可接受:
| 交互方式 | 体验 | 延迟感知 | 适用 |
|---|---|---|---|
| 全量等待 | 转圈加载直到完整输出 | 等待 = 完整生成时间 | 批处理、后台任务 |
| 流式输出 | 文字逐字呈现 | 等待 ≈ 首 token 延迟 | 对话、写作、搜索 |
| 半流式 | 分阶段逐步输出 | 等待 ≈ 关键中间产物 | Agent、工具调用 |
| 全双工 | 边输入边输出 | 双方同时进行 | 语音、实时协作 |
ℹ️ 核心洞察:流式改变的不仅是体验,更是系统语义——响应不再是“一个完成的结果”,而是一条持续到达的事件流,架构必须围绕事件流设计。
1.2 流式链路的分层
客户端(浏览器 / App / 终端)
│ EventSource / fetch(ReadableStream) / WebSocket
▼
API 网关(认证、限流、转发)
▼
流式服务(应用逻辑、Prompt 组装、工具调用)
▼
模型网关(多厂商路由、重试、熔断)
▼
模型推理(prefill + decode,逐 token 返回)
| 层 | 职责 | 流式相关挑战 |
|---|---|---|
| 客户端 | 渲染增量文本 | 解析流、增量渲染、中断按钮 |
| 网关 | 认证与限流 | 按 token 计量、连接保活 |
| 流式服务 | 业务编排 | 半流式、工具调用挂起 |
| 模型网关 | 多模型统一 | 厂商协议差异、超时重试 |
| 推理层 | 生成 token | TTFT、吞吐、decoding 策略 |
流式系统还要额外处理三个非流式没有的问题:连接生命周期(何时开始、何时结束)、部分输出(用户还没说完就被追问)、中断传播(用户点了停止后如何终止模型)。接口形态也因此从“返回完整字符串”变成“返回异步迭代器”:
async def answer_stream(query: str) -> AsyncIterator[str]:
"""流式接口:逐 token 产出,而非一次性返回完整结果。"""
async for token in model_stream(query):
yield token
二、流式协议与传输
2.1 SSE 与 WebSocket 的选型
| 维度 | SSE(Server-Sent Events) | WebSocket |
|---|---|---|
| 方向 | 服务端单向推送 | 双向全双工 |
| 协议 | 基于 HTTP,天然兼容 | 独立握手与帧协议 |
| 断线重连 | 内置自动重连 | 需自行实现 |
| 中间件 | 穿透 CDN/网关容易 | 网关需支持 Upgrade |
| 适用 | LLM 文本流、进度通知 | 语音流、实时协作、双向消息 |
| 复杂度 | 低 | 中高 |
一句话:纯 LLM 文本流优先 SSE——兼容性好、实现简单、自动重连;需要双向交互(打断、插话、连续对话)时再上 WebSocket。
2.2 SSE 的消息格式
SSE 用纯文本编码事件,字段以换行分隔:
event: token
data: {"delta": "你", "index": 0}
event: token
data: {"delta": "好", "index": 1}
event: done
data: {}
# sse_server.py — 用 FastAPI 实现 SSE 流式端点
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
import json
app = FastAPI()
async def sse_generator(query: str):
yield 'event: start\ndata: {"role":"assistant"}\n\n'
async for token in model_stream(query):
payload = json.dumps({"delta": token}, ensure_ascii=False)
yield f"event: token\ndata: {payload}\n\n"
yield "event: done\ndata: {}\n\n"
@app.get("/chat/stream")
async def chat_stream(query: str):
headers = {
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"X-Accel-Buffering": "no", # 禁止 Nginx 缓冲,保证逐条透传
}
return StreamingResponse(sse_generator(query),
media_type="text/event-stream",
headers=headers)
2.3 客户端消费 SSE
// sse_client.js — 浏览器端增量渲染
const source = new EventSource("/chat/stream?query=" + encodeURIComponent(q));
source.addEventListener("token", (e) => {
const data = JSON.parse(e.data);
render.appendDelta(data.delta); // 追加到光标处
autoScroll(); // 保持可见
});
source.addEventListener("done", () => {
source.close();
setDoneState();
});
source.onerror = () => {
// EventSource 默认自动重连;服务端需要保证幂等
console.log("连接断开,等待自动重连");
};
生产实现建议给每条流携带自增 index 与结束时的 finish_reason,方便客户端处理去重、乱序与结束判定——这是流式协议从“能跑”到“稳”的关键细节。
三、流式解析与渲染
3.1 增量解析:从字符流到结构化结果
流式到达的是字符串片段,但业务需要的是结构(JSON、Markdown、代码)。增量解析的关键是维持一个可恢复的状态机:
# incremental_parser.py — 增量 JSON 解析
import json
class IncrementalJsonParser:
def __init__(self):
self.buffer = ""
def feed(self, delta: str):
"""喂入新的字符片段,尝试把累积缓冲区解析为 JSON。"""
self.buffer += delta
try:
return json.loads(self.buffer), True
except json.JSONDecodeError:
return None, False # 不完整,等下一片
def is_done(self, finish_reason: str) -> bool:
return finish_reason == "stop"
3.2 增量渲染的防抖与光标稳定
流式渲染的经典坑是每次 delta 到达都整段重渲染,导致滚动位置与光标漂移。工程上用时间片批量渲染缓解:
| 策略 | 实现 | 效果 | 适用 |
|---|---|---|---|
| 逐 token 渲染 | 每个 delta 立即上屏 | 最丝滑、开销大 | 演示、对话 |
| 时间片批量 | 每 80ms 攒批渲染一次 | 平衡流畅与开销 | 生产默认 |
| 语义断句渲染 | 遇句号/换行才提交 | 更接近“打字”观感 | 长文生成 |
# render_scheduler.py — 时间片批渲染
import asyncio
class RenderScheduler:
def __init__(self, interval=0.08):
self.interval = interval
self.batch = []
async def submit(self, token: str):
self.batch.append(token)
async def flush_loop(self, on_flush):
while True:
await asyncio.sleep(self.interval)
if not self.batch:
continue
buf = "".join(self.batch)
self.batch.clear()
on_flush(buf)
四、半流式:把流式带进 Agent 与工具调用
4.1 什么是半流式
Agent 场景下,模型不是一次性生成答案,而是边思考、边调用工具、边组织回答。半流式指把这条“推理过程”分阶段推送给用户:先推 status 事件告知当前阶段(思考中 / 检索中),再推 tool_call 事件展示正在调用的工具,最后推 token 事件让最终回答逐字流出,以 done 结束。
4.2 工具调用流的暂停与恢复
流式中的函数调用很特殊:模型输出到 tool_call 时,流暂时中断等待工具执行,执行完后再恢复。实现上要用“事件 + 挂起任务”而非简单的 token 管道:
# tool_aware_stream.py — 流式工具调用编排
async def agent_stream(user_query: str):
"""Agent 流:遇到工具调用时暂停,执行后恢复。"""
pending_tool = None
async for event in model_agent_stream(user_query):
if event.type == "tool_call":
yield {"event": "tool_call", "name": event.name,
"args": event.arguments}
pending_tool = event
break
yield {"event": "token", "delta": event.delta}
if pending_tool:
result = await run_tool(pending_tool.name, pending_tool.arguments)
yield {"event": "tool_result", "result": result}
async for event in model_continue_stream(pending_tool, result):
yield {"event": "token", "delta": event.delta}
yield {"event": "done"}
五、首 token 延迟(TTFT)优化
5.1 TTFT 的构成
首 token 延迟(Time To First Token) 是流式体验的“第一印象”,由三部分组成:
| 环节 | 耗时构成 | 优化手段 |
|---|---|---|
| 网络往返 | 客户端→网关→服务 | 就近部署、连接复用 |
| Prefill | 模型处理整个输入(含历史) | 前缀缓存、压缩历史、预取 |
| 排队等待 | 推理集群负载 | 负载均衡、优先级、预留 |
def measure_ttft(start_ns: int, first_token_ns: int) -> float:
"""TTFT 度量:从请求发出到第一个 token 上屏(毫秒)。"""
return (first_token_ns - start_ns) / 1e6
5.2 前缀缓存降低 Prefill
流式对话的每次请求都携带全部历史,重复 prefill 是 TTFT 的最大来源。前缀缓存(Prompt Caching)让相同前缀命中缓存,prefill 几乎归零:
请求 1: [System][历史 H][问题 1] ← 全价
请求 2: [System][历史 H][问题 2] ← 前缀 [System][历史 H] 命中缓存
只对新问题 prefill,首 token 更快
# ttft_cache.py — 缓存感知的流式调用
def stream_with_cache(system_prompt, history, query):
"""相同前缀命中缓存,只对新 query 做 prefill。"""
messages = [{"role": "system", "content": system_prompt}] + history
messages.append({"role": "user", "content": query})
return client.chat.completions.create(
model="gpt-4o", messages=messages, stream=True)
ℹ️ 核心洞察:流式 + 前缀缓存是“快”与“省”的叠加——首 token 变快的同时,多轮对话重复输入的 token 成本可降 60% 以上。
TTFT 里的固定开销还包括 TLS 握手与连接建立,网关层可用连接池 + HTTP/2 多路复用(如 httpx 的 AsyncClient 复用 keepalive 连接)进一步压缩这部分耗时。
六、超时、中断与失败恢复
6.1 流式的超时矩阵
流式不是单次请求,超时要按阶段和空闲时长分别设置:
| 超时类型 | 含义 | 推荐值 |
|---|---|---|
| 连接建立 | 网关→模型建连 | 3-5s |
| TTFT 上限 | 首 token 迟迟不来 | 10-30s |
| 空闲间隔 | 相邻 token 间隔过长 | 5-10s |
| 总时长 | 整个流最长持续时间 | 120-300s |
6.2 用户中断:如何真正停止模型
用户点了“停止”,不是把前端清空就完事——下游模型还在继续烧 token。中断要逐级传播:
# cancellation.py — 中断的逐级传播
class StreamCancellation:
def __init__(self):
self._cancel = asyncio.Event()
async def cancel(self):
"""用户点击停止 → 置位取消事件。"""
self._cancel.set()
async def run(self, async_gen):
task = asyncio.create_task(self._drain(async_gen))
try:
await self._cancel.wait()
finally:
task.cancel()
await close_downstream() # 关闭 SSE 连接、释放资源
async def _drain(self, ag):
async for item in ag:
pass
网关侧同样要处理客户端断开:捕获 CancelledError 后调用厂商的 cancel 接口,避免下游继续烧 token:
async def stream_with_client_disconnect(query):
try:
async for token in model_stream(query):
yield token
except asyncio.CancelledError:
await model_client.cancel_last_request()
raise
6.3 断线重连与幂等续传
网络不可能永远稳定。断线重连有两个要求:位置可恢复(客户端知道渲染到哪)与下游幂等:
// resume_client.js — 基于 index 的续传
let lastIndex = 0;
function reconnectWithResume(query) {
const es = new EventSource(`/chat/stream?query=${query}&resume_from=${lastIndex}`);
es.addEventListener("token", (e) => {
const data = JSON.parse(e.data);
if (data.index <= lastIndex) return; // 丢弃重复片段
render.appendDelta(data.delta);
lastIndex = data.index;
});
}
七、流式可观测性
7.1 流式专属指标
传统 RPC 只有一次开始/结束,流式需要额外的过程指标:
| 指标 | 含义 | 告警阈值参考 |
|---|---|---|
| TTFT | 首 token 延迟 | P95 > 2s 告警 |
| Token 吞吐 | token/秒(解码速度) | 低于模型 50% 告警 |
| 完成率 | 完整走完 done 的连接占比 | < 95% 告警 |
| 中断率 | 用户主动取消占比 | 突增表示体验恶化 |
| 空闲超时 | token 间隔超时的流数 | 异常增长表示卡死 |
# stream_metrics.py — 流式指标埋点
def trace_stream_span(query, async_gen):
"""为每个流式请求建立 span,逐阶段打点。"""
span = tracer.start_span("llm.stream", attributes={"query": query})
span.set_attribute("ttft_ms", measure_ttft())
token_count = 0
async for item in async_gen:
token_count += 1
span.set_attribute("tokens_seen", token_count)
yield item
span.set_attribute("completed", got_done)
span.set_attribute("finish_reason", finish_reason)
span.end()
流式横跨 客户端→网关→服务→模型,trace id 要贯穿全程:用 X-Stream-Id 给每条流唯一标识,全链路日志、指标、采样评估都可以用它串起来。抽样评估时同样可以按 stream_id 把 token 按序重放成完整回答,沉淀进评测集——流式输出“流完但答错”也是一种质量回归。
八、流式架构的落地要点
8.1 网关:背压与限流
大量并发流会撑爆网关,需要背压(Backpressure):
# backpressure.py — 信号量限流 + 队列背压
import asyncio
class StreamGate:
def __init__(self, max_concurrent=200):
self.sem = asyncio.Semaphore(max_concurrent)
async def handle(self, query: str):
async with self.sem: # 超限排队而非丢弃
async for token in model_stream(query):
yield token
8.2 Nginx 反向代理的流式坑
Nginx 默认会缓冲响应,SSE 会被卡住直到结束。必须显式关闭缓冲:
# nginx.conf — 流式端点禁止缓冲
location /chat/stream {
proxy_pass http://llm_service;
proxy_buffering off; # 关闭缓冲
proxy_cache off;
proxy_read_timeout 300s; # 长连接不超时
chunked_transfer_encoding on;
}
8.3 流式的测试策略
流式逻辑不能只测“最终结果”,要测事件序列本身:
def test_stream_sequence():
"""断言事件顺序、done 结束、finish_reason 正确。"""
events = list(collect_events("今天天气怎么样?"))
kinds = [e["event"] for e in events]
assert kinds[0] == "start"
assert kinds[-1] == "done"
assert all(e["event"] != "done" for e in events[:-1])
text = "".join(e["delta"] for e in events if e["event"] == "token")
assert len(text) > 0
总结:流式工程的五条原则
| 原则 | 内容 |
|---|---|
| 1. 协议选型 | 纯文本流用 SSE,双向交互用 WebSocket |
| 2. 增量解析 | 用状态机解析片段,维护可恢复的缓冲区 |
| 3. 半流式编排 | 工具调用用事件 + 挂起,分阶段推送 |
| 4. 超时矩阵 | TTFT / 空闲 / 总时长分别设限,中断逐级传播 |
| 5. 全程可观测 | TTFT、吞吐、完成率、中断率全埋点 |
流式把 LLM 应用从“提交-等待-完成”的批处理模型,改造成“持续到达的事件流”模型。它的工程重心因此从“怎么生成”转向“怎么传输、怎么增量渲染、怎么中断、怎么度量”。掌握协议选型、增量解析、半流式编排与超时中断治理,你就能让每一次生成都成为一次流畅、可控、可观测的实时体验。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。