Skip to content

第17章 · 全栈流式交付

live 转发路径在本仓库证据中标记 unverified-live,默认证据线是离线录制流回放。

先修:第14章(真实 LLM API、SSE chunk 结构、parse_sse_events 行缓冲解析器)。本章是它的前端延伸:第14章教你"流怎么解析",本章教你"流怎么从服务端一直送到用户眼睛里"。

本章目标

  • 说清 SSE 帧格式(data: 帧、\n\n 事件边界、: 心跳注释、[DONE] 哨兵),并能用 stdlib 写一个转发 LLM 流的 SSE 端点。
  • 口述 OpenAI chat.completion.chunk 与 Gemini :streamGenerateContent 的结构对照:增量字段、首/末 chunk 语义、哨兵与 usage 的差异。
  • 用 fetch + ReadableStream(而非 EventSource)消费 LLM 流:行缓冲、流式 TextDecoder、AbortController 中断,并解释为什么 LLM 流没有断点续传。
  • 亲手制造两类流式故障并修复:把 TCP chunk 当消息边界(半行 JSON 崩解析)、逐 chunk 非流式 decode(多字节字符乱码)。
  • 实现 partial JSON 容错解析器,把 tool-call 式的结构化工件分片流渲染成"字段随生成逐个填入"的实时表单。
  • 说清 Markdown 增量渲染的核心矛盾(任意前缀不一定是合法 Markdown)与工程做法(安全边界缓冲 / 容错解析器),并识别逐 chunk 全量 re-parse 的 O(n²) 反模式。

接口 shape 速览

对象shape/结构约束
SSE 帧data: <payload>\n\n: 开头为心跳注释事件边界是空行,不是 TCP chunk;一个 JSON 事件可被拆进多个 chunk
OpenAI 流 chunkchoices[0].delta.content(增量);首 chunk delta.role;末 chunk finish_reason结束哨兵 data: [DONE];usage 需 stream_options.include_usage;tool call 参数是分片字符串需拼接
Gemini 流 chunkcandidates[0].content.parts[0].text(增量)[DONE] 哨兵,以 finishReason: "STOP" 结束;usageMetadata 每个 chunk 都有
前端消费fetch POST → res.bodyTextDecoderStream/增量 decoder → 行缓冲 → 帧解析不用 EventSource(GET-only、无 body、自动重连≠断点续传)
中断AbortController.abort() → reader 取消 → 服务端写失败关连接重连 = 从头重新生成,无续传语义
结构化工件JSON 字符串分片,任意前缀大多不是合法 JSONpartial JSON 容错解析(自动闭合 + 前缀回退),字段逐个填入 UI

从零实践(三个模块)

服务端参考实现是 python/labs/streaming_server.py(stdlib-only,http.server):本地假流发生器按固定间隔吐录制 chunk,完全离线可测;--live 可选路径复用第14章的 GeminiRoundtripClient.stream_chat 做真实转发(无 key 时标 unverified-live)。前端演示组件是 docs/.vitepress/theme/components/StreamingDemo.vue,已注册进主题并嵌入下方 B2 节。

B1 · 服务端 SSE 与 chunk 协议

SSE 的全部协议就四条:Content-Type: text/event-stream;事件以 data: <payload>\n\n 为单位;: 开头的行是心跳/注释(防代理超时断连);空行是事件边界。唯一容易翻车的是第四条推论:TCP 不保证任何边界——一个 JSON 事件可能被拆进多个 chunk,也可能多个事件挤在一个 chunk 里。所以消费端必须按行缓冲(第14章的 parse_sse_events 已经实现:残留半行留到下次、\r\n 处理、多行 data: 拼接、忽略注释),把字节 chunk 当消息边界是本模块的头号故障。

结构对照(极好的教学对照组,面试高频):

关注点OpenAI 兼容流Gemini 原生流
增量文本choices[0].delta.contentcandidates[0].content.parts[0].text
首 chunkdelta.role 建立角色无特殊角色 chunk
结束信号data: [DONE] 哨兵 + 末 chunk finish_reason无哨兵finishReason: "STOP"
usage默认不带;stream_options.include_usage=true 时末 chunk 带usageMetadata 每个 chunk 都有
tool call参数是 JSON 字符串分片,按 index 拼接后解析一次functionCall part(非流式分片语义不同)

动手:

  1. 启动本地假流端点并用 curl -N 裸看原始帧(-N 关闭缓冲,逐字节看 SSE 帧到达):

    bash
    PYTHONPATH=python .venv/bin/python -m labs.streaming_server --port 8765 &
    curl -N http://127.0.0.1:8765/stream        # chat 文本流
    curl -N http://127.0.0.1:8765/stream-json   # 结构化 JSON 分片流
  2. 观察服务器以 8 字节固定边界切块:帧里的中文必然被切在 UTF-8 序列中间——这正是 B2 要消化的现实。

  3. 改写 --chunk-size 为 1,确认解析结果不变(行缓冲 + 增量 decode 对任意切块不变式)。

  4. 有 key 时自选:--live 转发真实 Gemini 流,对照录制流验证帧形状一致(无 key 标 unverified-live,不要编造帧内容)。

B2 · 前端消费与中断

用 fetch + ReadableStream,不用 EventSource。 原因两条:EventSource 只支持 GET(LLM API 全是 POST + JSON body + 鉴权 header);它的自动重连在 LLM 场景是缺点而非优点——重连不等于断点续传,服务端没有任何游标,重连只会从头重新生成一遍(还要再烧一遍 token)。

最小消费骨架(Vue/TS,前端部分只给骨架,难点在 LLM 特有语义):

ts
const ctrl = new AbortController()
const res = await fetch(url, { method: 'POST', headers, body, signal: ctrl.signal })
const reader = res.body!.pipeThrough(new TextDecoderStream()).getReader()
let buf = ''
for (;;) {
  const { done, value } = await reader.read()
  if (done) break
  buf += value
  const lines = buf.split('\n')
  buf = lines.pop()!           // 残留半行留到下次
  for (const line of lines) { /* 解析 data: 帧;[DONE] 终止;ref 累加触发响应式 */ }
}
// "停止生成"按钮 = ctrl.abort();中断后没有续传,只有重新生成

TextDecoder 流式模式是第二个 LLM 特有坑:token 边界可能切在多字节 UTF-8 字符中间(CJK 全是 3 字节序列)。逐 chunk 独立 new TextDecoder().decode(chunk) 会把被切断的字节解码成 U+FFFD 乱码;TextDecoderStreamdecoder.decode(chunk, { stream: true }) 会持有未完成的字节序列直到后续 chunk 补齐。

交互:流式消费管线(行缓冲 × 增量 decode × partial JSON)

真实端点(可选,POST;留空则跑模拟流)

真实路径用 fetch + ReadableStream + AbortController,不用 EventSource(只支持 GET、无 body、自动重连≠断点续传)。跨域需服务端开 CORS。

状态: 待连接 | 事件: 0 | 心跳: 0 | [DONE]: 否

流式 decode(正确)
逐 chunk 非流式 decode(错误示范:乱码对照)
解析后累加的 delta.content
(空)

动手:在 StreamingDemo 组件里把"模拟 chunk 字节边界"滑块拖到 1–8,对照左右两个面板——左边流式 decode 完好,右边非流式 decode 满屏 。然后回答:为什么把 chunk 调大到 64 乱码也不保证消失?(边界由网络栈决定,不由你决定;任何固定假设都会被打破。)

B3 · 增量渲染进阶

Markdown 增量渲染的核心矛盾:流的任意前缀不一定是合法 Markdown——代码 fence、表格、链接都可能未闭合。工程做法两档:缓冲到安全边界(如 \n\n 段落边界)再交给渲染器;或用容错增量解析器(markdown-it 这类对未闭合结构天然宽容的解析器属于后者)。逐 chunk 全量 re-parse 是 O(n²) 反模式:n 个 token 各触发一次全长解析,长回答下 CPU 曲线会直接教做人。

结构化工件流式渲染(学习者是前端大牛,这是炫技空间,难点全在 LLM 分片语义):tool call 参数是 JSON 字符串分片,任意前缀大多不是合法 JSON——但用户想看到"字段随生成逐个填入"的实时表单。解法是 partial JSON 容错解析:扫描括号栈与字符串状态,自动闭合未完成的字符串/对象/数组,解析失败就回退到最后一个完整键值对再闭合。StreamingDemo 的 JSON 场景内置了一个参考实现(parsePartialJson,约 40 行)。

动手:

  1. 让(假)模型流式输出固定 schema JSON:curl -N http://127.0.0.1:8765/stream-json,确认几乎每个分片单独 JSON.parse 都失败,拼接后才是完整对象(测试 test_json_fragments_only_concatenate_to_valid_json 断言了这一点)。
  2. 在组件里切到"结构化 JSON 流"场景,观察表单字段逐个填入;把字节边界拖到 1,确认填入结果不变。
  3. 进阶自命题:给 partial JSON 解析器加"数组元素逐个出现"的动画;或写一个 Markdown 安全边界缓冲器(只在 \n\n 处切渲染)。

故障注入与预期信号

注入预期失败信号修复后证据
把 TCP chunk 当消息边界,逐 chunk split 后不缓冲跨 chunk 的半行 JSON 触发解析异常,流随机崩行缓冲(残留半行留下次)通过半行 fixture;test_half_line_reassembly_matches_recorded_events 绿
逐 chunk new TextDecoder().decode()(非流式模式)多字节字符被切成 U+FFFD 乱码(流式 decode({stream:true}/TextDecoderStream)逐字节还原;test_incremental_decode_recovers_exact_text 绿、test_naive_decode_is_mojibake_or_raises 证明乱码存在
用 EventSource 消费 LLM APIPOST body/鉴权 header 无处安放;断线后自动重连导致从头重新生成、重复扣额度改 fetch + ReadableStream + AbortController;中断语义明确为"终止",无幽灵重连
每个 delta 都 JSON.parse 工具参数参数分片不是合法 JSON,解析异常index 拼接全部片段后只解析一次(第14章 accumulate_chat_chunks 同款语义)
对 partial JSON 直接 JSON.parse 驱动表单流式期间全程抛异常,表单只在最后一刻闪现容错解析(自动闭合+前缀回退)后字段随生成逐个填入;组件 JSON 场景可复现
逐 chunk 全量 re-parse Markdown解析耗时随长度平方增长,长回答掉帧安全边界(\n\n)缓冲或容错增量解析器,渲染次数与段落数成正比
未闭合代码 fence 直接交给渲染器后半段回答被吞进代码块、样式雪崩边界缓冲后渲染,或渲染前临时补闭合 fence
中断后假设"续传",用同一连接重试服务端无游标,重复生成 + 重复计费AbortController 终止即终态;重发=新请求,UI 明确标注"已中断"

规范与延伸

前端/Agent 迁移

这一章几乎全是你的主场,迁移价值在语义层:SSE 行缓冲就是你在 WebSocket 分帧里做过的消息重组;AbortController 就是"停止生成"按钮;partial JSON 容错解析就是把 tool call 分片当成"渐进式表单协议"。反向迁移一条:模型流是 data 不是 instruction——流式渲染出的 Markdown/JSON 只进 UI,不因为"它能流"就获得执行权;tool call 拼完整后仍走第13章的 schema/policy 校验。Agent 侧的同构物是第13章的 checkpoint:流式输出本身不可恢复(无游标),可恢复的是 checkpoint 里的消息历史——中断后的正确动作是从 checkpoint 重放,不是"接着收"。

口述与自测(不看资料,5–10 分钟)

  • 画出从 TCP 字节到屏幕上字符的完整管线:chunk → 增量 decode → 行缓冲 → SSE 帧 → delta 累加 → 渲染;指出半行、心跳、[DONE]、多字节切断各在哪一步被处理。
  • 背出 OpenAI 与 Gemini 流式结构的四个差异(增量字段、哨兵、usage、结束信号),并解释为什么"无哨兵"对客户端意味着什么。
  • 解释为什么 EventSource 不适合 LLM 流(两条理由),以及"自动重连"为什么会变成重复扣费。
  • 口述 partial JSON 容错解析的算法(括号栈 + 字符串状态 + 前缀回退),并说明它为什么不能用于"校验",只能用于"展示"。
  • 说明 Markdown 增量渲染的矛盾与两档工程解法,背出 O(n²) 反模式的成因。

动手实验

把"全栈流式交付"拆成离线可验证的确定性证据:帧格式、半行重组、多字节切断、中断语义全部由录制流回放覆盖;live 转发只在你自己的 key 下单独核验并标记。

环境准备

bash
cd <仓库>
export PYTHONPATH="$PWD/python"

命令与预期输出

bash
# 1) 离线回放测试(CI 默认线,无需 key、无网络)
.venv/bin/python -m pytest python/tests/test_streaming_server.py -q

# 2) 启动假流端点,curl -N 裸看原始帧
PYTHONPATH=python .venv/bin/python -m labs.streaming_server --port 8765 &
curl -N http://127.0.0.1:8765/stream | head -c 400
text
...........                                                              [100%]
11 passed in 2.10s

: keep-alive

data: {"id": "chatcmpl-recorded", "object": "chat.completion.chunk", "choices": [{"index": 0, "delta": {"role": "assistant"}, "finish_reason": null}]}

data: {"id": "chatcmpl-recorded", "object": "chat.completion.chunk", "choices": [{"index": 0, "delta": {"content": "流式交付的关键:"}, "finish_reason": null}]}

判定信号:
  帧格式断言通过(data: 前缀、\n\n 边界、: 心跳、[DONE] 哨兵)
  8 字节固定切块下:行缓冲重组后事件与录制源逐一相等;流式 decode 逐字节还原;
  非流式 decode 产生 U+FFFD 乱码;客户端中断后服务端观察到连接关闭

live 核验(有 key 时自选,不计入本仓库证据):

bash
export GEMINI_API_KEY="<your-key>"
PYTHONPATH=python .venv/bin/python -m labs.streaming_server --port 8765 --live &
curl -N http://127.0.0.1:8765/stream

概念图

资源 / 成本 / 隐私

离线回放路径 gross cost 为 0、无网络。live 转发走免费层时账单为零但额度有限(同第14章纪律:RPM/TPM/RPD 随账户变动,不硬编码数字);免费层请求数据会被 Google 用于改进产品,敏感内容一律不发;key 只走环境变量,不落盘、不进 fixture、不进帧、不进日志。录制流内容全部为课程原创合成文本。

Evidence

仓库当前机器证据(只读快照)

evidence/17-streaming-v1.json 是当前 checkout 的脱敏机器运行记录,只覆盖离线录制流回放:test_streaming_server.py 11 项全部通过(帧格式、半行重组、多字节切断、中断观察)。--live 转发路径在无 key 环境下标记 unverified-live 并写入 known_failures。本模块已登记进 evidence/module-manifest-v1.json。前端组件 StreamingDemo.vue 已注册进主题,其演示行为不在本证据覆盖范围内。

学习者提交模板(待填写,不是当前机器证据)

复制下面模板并填写自己的真实运行结果。所有 <...> 都是未填写状态;actualartifacts 尤其不能被当作已运行或已通过。artifacts 必须替换为本次提交中真实存在的仓库相对路径。live 指标(若有)必须来自你自己的 key 的真实调用,并在 known_failures 注明数据边界。

yaml
schema: learn-llm.evidence.v1
module: 17-streaming
commit: <learner-commit-sha>
verified_at: <iso-date>
environment: <sanitized-python-device>
seed: 10
commands:
  - PYTHONPATH=python python -m pytest python/tests/test_streaming_server.py -q
metrics:
  - name: offline_streaming_tests_passed
    expected: 11
    actual: <recorded-value>
  - name: sse_frames_reassembled_equal_source
    expected: true
    actual: <recorded-value>
  - name: multibyte_severed_chunks_detected
    expected: ">0"
    actual: <recorded-value>
  - name: abort_observed_server_side
    expected: true
    actual: <recorded-value>
  - name: live_forward_calls
    expected: unverified-live
    actual: <recorded-value-or-unverified-live>
artifacts:
  - <learner-repo-relative-artifact-path>
cost:
  gross_usd: 0
  credit_usd: 0
licenses:
  - source: <source>
    version: <version>
    license: <license>
    attribution: <attribution>
    redistribution: <redistribution>
known_failures:
  - <sanitized-failure-or-none>

只有页面勾选、没有流式解析与中断断言的测试证据时,本章保持 gate

下一步

私有学习站 · 原理从零构建 · 勿提交个人隐私或密钥