LLM 流式输出与实时交互工程

人类对“等待”的耐心是有限的:当用户发出一个问题,超过 3 秒没有反馈,流失率就会显著上升。LLM 应用因此把“生成中”的状态转化为体验的一部分——流式输出让模型一个字一个字地“打字”,用户感受到的不是等待,而是思考的进行。但流式不是简单地把 HTTP 响应切成小块。它牵扯协议选型、增量解析、中断语义、首 token …

人类对“等待”的耐心是有限的:当用户发出一个问题,超过 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 计量、连接保活
流式服务业务编排半流式、工具调用挂起
模型网关多模型统一厂商协议差异、超时重试
推理层生成 tokenTTFT、吞吐、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 应用从“提交-等待-完成”的批处理模型,改造成“持续到达的事件流”模型。它的工程重心因此从“怎么生成”转向“怎么传输、怎么增量渲染、怎么中断、怎么度量”。掌握协议选型、增量解析、半流式编排与超时中断治理,你就能让每一次生成都成为一次流畅、可控、可观测的实时体验。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「LLM」更多文章

  1. 推理增强技术工程化:CoT/ToT/ReAct
  2. 模型路由与选型:大小模型分层调度
  3. 混合检索 RAG:BM25 与向量融合