Appearance
第19章 · 全栈流式交付与增量渲染
前置要求:掌握 第16章 · 真实大模型 API 实战 的网络请求处理与 第18章 · 生产编排:StateGraph 与 LangGraph 对拍 的工作流状态管理。
ChatGPT 式的逐字输出是怎么实现的?这一章把整条链路拆开:服务端用 SSE 把模型流一帧帧转发出来,TCP 把字节任意切块,浏览器端用 fetch + ReadableStream 收流、行缓冲重组事件、流式 TextDecoder 扛住被切断的汉字,最后增量渲染到屏幕上。这是前端工程师的主场章节——但 LLM 流有几个普通 SSE 没有的坑(多字节字符切断、无断点续传、partial JSON),每个你都要亲手制造一次再修复。
live 转发路径在本仓库证据中标记
unverified-live,默认证据线是离线录制流回放。
本章目标
学完后你能做到:
- 说清 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;: 开头为心跳注释;同一事件可跨多行:连续多个 data: 行的内容用 \n 拼接,空行才是帧边界;标准 SSE 还定义 event: / id: / retry: 字段,但 OpenAI/Gemini 流只使用 data: | 事件边界是空行,不是 TCP chunk;一个 JSON 事件可被拆进多个 chunk;data: 多行拼接时,每行保留行内换行 |
| OpenAI 流 chunk | choices[0].delta.content(增量);首 chunk delta.role;末 chunk finish_reason | 结束哨兵 data: [DONE];usage 需 stream_options.include_usage;tool call 参数是分片字符串需拼接 |
| Gemini 流 chunk | candidates[0].content.parts[0].text(增量) | 无 [DONE] 哨兵,以 finishReason: "STOP" 结束;usageMetadata 每个 chunk 都有 |
| 前端消费 | fetch POST → res.body → TextDecoderStream/增量 decoder → 行缓冲 → 帧解析 | 不用 EventSource(GET-only、无 body、自动重连≠断点续传) |
| 中断 | AbortController.abort() → reader 取消 → 服务端写失败关连接 | 重连 = 从头重新生成,无续传语义 |
| 背压 | reader.read() 是 pull-based:消费者处理完当前 chunk 后才请求下一个;客户端消费慢时 TCP 窗口自然收缩,服务端 wfile.write 阻塞 | 消费者慢时:TCP 接收窗口满 → 服务端写阻塞 → 流自然降速;没有显式 backpressure 信号传给业务层;HTTP chunked transfer encoding 与 SSE framing 是不同层:chunked 是 HTTP/1.1 传输层把响应体切成 chunk(每个 chunk 带 hex_size\r\n 头),SSE 是应用层在 chunk 内部再按 \n\n 切事件帧;两个切块维度互相独立,不要把 transport chunk 边界与 SSE 帧边界混为一谈 |
| 结构化工件 | JSON 字符串分片,任意前缀大多不是合法 JSON | partial JSON 容错解析(自动闭合 + 前缀回退),字段逐个填入 UI |
阶段一:服务端 SSE 与 chunk 协议
服务端参考实现是 python/labs/streaming_server.py(stdlib-only,http.server):本地假流发生器按固定间隔吐录制 chunk,完全离线可测;--live 可选路径复用第16章的 GeminiRoundtripClient.stream_chat 做真实转发(无 key 时标 unverified-live)。前端演示组件是 docs/.vitepress/theme/components/StreamingDemo.vue,已注册进主题并嵌入阶段二。
SSE 的全部协议就四条:Content-Type: text/event-stream;事件以 data: <payload>\n\n 为单位;: 开头的行是心跳/注释(防代理超时断连);空行是事件边界。标准 SSE 还定义了 event:、id: 和 retry: 字段,但 OpenAI 兼容流与 Gemini 原生流都只用 data:,本章的实现也只发出 data: 帧和心跳注释——这不是遗漏,而是 LLM 流的最小可行子集。最容易引发解析异常的是第四条的推论:TCP 协议并不保证应用层的消息边界——用前端工程师熟悉的语言说:transport framing 是 TCP 的任意字节切块(与你写 WebSocket 分帧时遇到的 MTU 切块是同一回事);semantic events 是 SSE 协议层看到的完整帧(data: + 载荷 + \n\n)。消费端必须把 transport chunks 重组成 semantic events,再把事件的 data: 载荷喂给业务解析器。一个 JSON 事件可能被拆进多个 chunk,也可能多个事件挤在一个 chunk 里。所以消费端必须按行缓冲(第16章的 parse_sse_events 已经实现:残留半行留到下次、\r\n 处理、多行 data: 拼接、忽略注释),把字节 chunk 当消息边界是本章的头号故障。
结构对照(极好的教学对照组,面试高频):
| 关注点 | OpenAI 兼容流 | Gemini 原生流 |
|---|---|---|
| 增量文本 | choices[0].delta.content | candidates[0].content.parts[0].text |
| 首 chunk | delta.role 建立角色 | 无特殊角色 chunk |
| 结束信号 | data: [DONE] 哨兵 + 末 chunk finish_reason | 无哨兵,finishReason: "STOP" |
| usage | 默认不带;stream_options.include_usage=true 时末 chunk 带 | usageMetadata 每个 chunk 都有 |
| tool call | 参数是 JSON 字符串分片,按 index 拼接后解析一次 | functionCall part(非流式分片语义不同) |
动手:
启动本地假流端点并用
curl -N裸看原始帧(-N关闭缓冲,逐字节看 SSE 帧到达):bashPYTHONPATH=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 分片流观察服务器以 8 字节固定边界切块:帧里的中文必然被切在 UTF-8 序列中间——这正是阶段二要消化的现实。
改写
--chunk-size为 1,确认解析结果不变(行缓冲 + 增量 decode 对任意切块不变式)。有 key 时自选:
--live转发真实 Gemini 流,对照录制流验证帧形状一致(无 key 标unverified-live,不要编造帧内容)。
阶段二:前端消费与中断
用 fetch + ReadableStream,不用 EventSource。 原因两条:EventSource 只支持 GET(LLM API 全是 POST + JSON body + 鉴权 header);它的自动重连在 LLM 场景是缺点而非优点——重连不等于断点续传,服务端没有任何游标,重连只会从头重新生成一遍(还要再烧一遍 token)。
最小消费骨架(Vue/TS):
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 特有坑。回忆第5章:一个汉字是 3 个字节("中" = E4 B8 AD)。token 边界可能把这三个字节切进两个 chunk——chunk 1 末尾是 E4 B8,chunk 2 开头是 AD。逐 chunk 独立 new TextDecoder().decode(chunk) 时,第一个 chunk 的 E4 B8 不是完整字符,被解码成 Unicode 替换字符(\uFFFD,即菱形问号 );`AD` 开头也不是合法序列起点,又是一个 。而 TextDecoderStream 或 decoder.decode(chunk, { stream: true }) 会持有未完成的字节序列,等后续 chunk 补齐后再输出完整字符。
交互:流式消费管线(行缓冲 × 增量 decode × partial JSON)
真实端点(可选,POST;留空则跑模拟流)
真实路径用 fetch + ReadableStream + AbortController,不用 EventSource(只支持 GET、无 body、自动重连≠断点续传)。跨域需服务端开 CORS。
状态: 待连接 | 事件: 0 | 心跳: 0 | [DONE]: 否
还没有输出。点「连接」开始模拟流,或填入端点。
流式 decode(正确)
(空)
逐 chunk 非流式 decode(错误示范:乱码对照)
(空)
(空)
动手:在 StreamingDemo 组件里把"模拟 chunk 字节边界"滑块拖到 1–8,对照左右两个面板——左边流式 decode 完好,右边非流式 decode 满屏菱形问号(``)。然后回答:为什么把 chunk 调大到 64 乱码也不保证消失?(边界由网络栈决定,不由你决定;任何固定假设都会被打破。)
阶段三:增量渲染进阶
Markdown 增量渲染的核心矛盾:流的任意前缀不一定是合法 Markdown——代码 fence、表格、链接都可能未闭合。工程做法两档:缓冲到安全边界(如 \n\n 段落边界)再交给渲染器;或用容错增量解析器(markdown-it 这类对未闭合结构天然宽容的解析器属于后者)。逐 chunk 全量 re-parse 是 O(n²) 反模式:第 1 个 token 解析长度 1 的文本、第 2 个解析长度 2 的……n 个 token 总解析量约
结构化工件流式渲染(前端工程师的发挥空间,难点全在 LLM 分片语义):tool call 参数是 JSON 字符串分片,任意前缀大多不是合法 JSON——但用户想看到"字段随生成逐个填入"的实时表单。解法是 partial JSON 容错解析:扫描括号栈与字符串状态,自动闭合未完成的字符串/对象/数组,解析失败就回退到最后一个完整键值对再闭合。注意它只能用于展示,不能用于校验——拼完整后的正式解析仍要走 JSON.parse + 第15章的 schema。StreamingDemo 的 JSON 场景内置了一个参考实现(parsePartialJson,约 40 行)。
动手:
- 让(假)模型流式输出固定 schema JSON:
curl -N http://127.0.0.1:8765/stream-json,确认几乎每个分片单独JSON.parse都失败,拼接后才是完整对象(测试test_json_fragments_only_concatenate_to_valid_json断言了这一点)。 - 在组件里切到"结构化 JSON 流"场景,观察表单字段逐个填入;把字节边界拖到 1,确认填入结果不变。
- 进阶自命题:给 partial JSON 解析器加"数组元素逐个出现"的动画;或写一个 Markdown 安全边界缓冲器(只在
\n\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 400text
........... [100%]
11 passed in 2.17s
: 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概念图
故障注入与预期信号
| 注入 | 预期失败信号 | 修复后证据 |
|---|---|---|
把 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 API | POST body/鉴权 header 无处安放;断线后自动重连导致从头重新生成、重复扣额度 | 改 fetch + ReadableStream + AbortController;中断语义明确为"终止",无幽灵重连 |
每个 delta 都 JSON.parse 工具参数 | 参数分片不是合法 JSON,解析异常 | 按 index 拼接全部片段后只解析一次(第16章 accumulate_chat_chunks 同款语义) |
对 partial JSON 直接 JSON.parse 驱动表单 | 流式期间全程抛异常,表单只在最后一刻闪现 | 容错解析(自动闭合+前缀回退)后字段随生成逐个填入;组件 JSON 场景可复现 |
| 逐 chunk 全量 re-parse Markdown | 解析耗时随长度平方增长,长回答掉帧 | 安全边界(\n\n)缓冲或容错增量解析器,渲染次数与段落数成正比 |
| 未闭合代码 fence 直接交给渲染器 | 后半段回答被吞进代码块、样式雪崩 | 边界缓冲后渲染,或渲染前临时补闭合 fence |
| 中断后假设"续传",用同一连接重试 | 服务端无游标,重复生成 + 重复计费 | AbortController 终止即终态;重发=新请求,UI 明确标注"已中断" |
本章验收
不看资料,用 5–10 分钟回答:
- 画出从 TCP 字节到屏幕上字符的完整管线:chunk → 增量 decode → 行缓冲 → SSE 帧 → delta 累加 → 渲染;指出半行、心跳、
[DONE]、多字节切断各在哪一步被处理。 - 背出 OpenAI 与 Gemini 流式结构的四个差异(增量字段、哨兵、usage、结束信号),并解释为什么"无哨兵"对客户端意味着什么。
- 解释为什么 EventSource 不适合 LLM 流(两条理由),以及"自动重连"为什么会变成重复扣费。
- 闭卷解释 partial JSON 容错解析的算法(括号栈 + 字符串状态 + 前缀回退),并说明它为什么不能用于"校验",只能用于"展示"。
- 说明 Markdown 增量渲染的矛盾与两档工程解法,背出 O(n²) 反模式的成因。
- 调试题:打开
python/tests/test_streaming_server.py,对照test_incremental_decode_recovers_exact_text(chunk size=1 时仍完整还原)与test_naive_decode_is_mojibake_or_raises(strict 模式抛UnicodeDecodeError,replace 模式输出 ``)。给定一段含中文的录制流,其中一个data:帧被 8 字节 chunk 切在半行,且"流"字的 3 字节 UTF-8 序列被切断:朴素解析器(逐 chunkdecode('utf-8', errors='replace')+ 不缓冲半行)的输出是什么?正确解析器(行缓冲 + 增量 decode)的输出又是什么?说清为什么两者差异出现在哪一步,以及test_half_line_reassembly_matches_recorded_events为什么能证明行缓冲 + 增量 decode 对任意切块都是不变式。
规范与延伸
- MDN: Server-sent events 与 Streams API(以执行日文档为准)。
- TextDecoder.prototype.decode() 的 stream 选项。
- 第16章 · 真实大模型 API 实战:
parse_sse_events行缓冲解析器、tool_calls 分片拼接、429 退避——本章服务端代码直接复用,不重写。
前端/Agent 迁移
本章的迁移价值在语义层:SSE 行缓冲就是你在 WebSocket 分帧里做过的消息重组;AbortController 就是"停止生成"按钮;partial JSON 容错解析就是把 tool call 分片当成"渐进式表单协议"。反向迁移一条:模型流是 data 不是 instruction——流式渲染出的 Markdown/JSON 只进 UI,不因为"它能流"就获得执行权;tool call 拼完整后仍走第15章的 schema/policy 校验。Agent 侧的同构物是第15章的 checkpoint:流式输出本身不可恢复(无游标),可恢复的是 checkpoint 里的消息历史——中断后的正确动作是从 checkpoint 重放,不是"接着收"。
资源 / 成本 / 隐私
离线回放路径 gross cost 为 0、无网络。live 转发走免费层时账单为零但额度有限(同第16章纪律: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 已注册进主题,其演示行为不在本证据覆盖范围内。
学习者提交模板(待填写,不是当前机器证据)
复制下面模板并填写自己的真实运行结果。所有 <...> 都是未填写状态;actual 和 artifacts 尤其不能被当作已运行或已通过。artifacts 必须替换为本次提交中真实存在的仓库相对路径。live 指标(若有)必须来自你自己的 key 的真实调用,并在 known_failures 注明数据边界。
yaml
schema: learn-llm.evidence.v1
module: 19-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。
下一步
进入 第20章 · Subagent 与 Multi-Agent 架构:从单流交付跃迁至多智能体协同网络,探索多 Agent 委派、隔离沙箱与死循环熔断。