3.3 流式输出(Streaming)
本章解决什么问题:上一章我们把「一次对话」拆成了消息与 token。但真实的模型回复不是「一瞬间出现的一整段文字」——它是一个字一个字生成出来的,慢的时候要几十秒。如果程序非要等它写完才显示,用户就得盯着空白屏幕干等。本章讲清楚流式输出(Streaming)是什么形状:它不是「一串字符」,而是一串事件;讲清楚这串事件怎么从模型服务器一路送到你的终端,以及它为什么正好长成 2.6 学过的异步迭代器(Async Iterator)的样子。
前置知识:3.1 LLM 与模型 API、3.2 消息、上下文与 token(知道一次请求发出去什么、收回来什么);2.6 异步迭代器与 for await(
async function*与for await..of);2.3 union、字面量与可辨识联合(事件类型是一个可辨识联合)。学习目标:读完本章后你能
- 说出流式输出到底改善了什么、又没有改善什么,并举出两个「不流式就做不到」的功能;
- 描述一次回复的事件序列形状:
start→ 若干delta→end,并解释为什么中间不能只发纯文本;- 用一句话解释 SSE(Server-Sent Events,服务器发送事件)是什么,以及它在 Pi 的数据链路里处于哪一层;
- 用
async function*写出一个逐块吐事件的假模型,用for await..of消费出打字机效果,并把事件汇总回一条完整消息;- 说清楚「部分 JSON」问题是怎么产生的,为什么工具调用的参数在流式期间不能直接
JSON.parse。
建立直觉:30 秒白屏 vs 立刻看到第一个词
假设一个大语言模型(LLM,Large Language Model)每秒生成 30 个 token(这是为了算账随手取的数量级,不同模型差别很大),你让它写一段 600 token 的解释——生成完要 20 秒。摆在程序面前的是两种拿结果的方式:
- 非流式:发出请求,等 20 秒,一次性收到整段文字。这 20 秒里屏幕上什么都没有,用户无法区分「模型在思考」和「程序已经卡死」,只能盯着光标。
- 流式:模型每生成一小块就立刻发出来,程序收到一块就显示一块。第一块大约 0.5 秒后到达,文字随后一点点长出来,20 秒后同样收完。
关键要说清楚:流式一点也没有让模型生成得更快。总时间还是 20 秒,token 还是那么多,费用一分不少。它改变的只有一件事——你什么时候能看到第一个词。但这一件事牵出了三个非流式做不到的功能:
- 可感知的响应速度:20 秒的空白被摊成了几百次「又长出几个字」。人对「在动」和「不动」的容忍度差着数量级。
- 可中途叫停:看到模型第二句话就跑偏了,用户按下 Esc,请求当场取消——省下剩下 18 秒和相应的 token 费用。非流式模式下你连「跑偏了」都不知道,只能等它写完。(怎么取消是 2.7 讲过的 AbortSignal,在 Pi 里怎么落地见 5.6。)
- 边到边处理:界面可以边收边渲染 Markdown;Agent 可以在工具调用的参数还没收完时就先把「正在读取文件 xxx」显示出来。数据不必攒齐才能用。
代价也要摆出来:错误可能发生在「已经显示了一半」的时候,协议必须能表达「中途失败」;消费者也从「拿到一个结果」变成「处理一串事件、最后自己拼出结果」。这是那三条好处的账单。
流式的形状:一串事件,不是一串字符
最容易想到的流式形状是「一串字符串片段」:"流式"、"输出的"、"意思是"……收到就往屏幕上追加。这个模型在只有纯文本回复时够用,但一遇到真实场景就散架:
- 模型这一轮可能既说了话,又要求调用一个工具(3.4 的主题),还可能夹着一段「思考」内容。纯字符串流没法表达「这一块属于哪部分」。
- 消费者需要知道边界:文本块从哪开始、到哪结束——不然没法决定何时给这段文字收尾、何时开始渲染下一块。
- 这一轮结束时还有一批只有终局才知道的信息:为什么停下来(说完了?被截断了?要调工具?)、用了多少 token、以及权威的完整消息。
所以真实的流式协议把结果表达成一串带类型的事件。最小的一套是三类:
start:这一轮回复开始了(此时还一个字都没有);- 若干
delta:每个带一小块增量内容(delta就是「增量」); end/done:结束,并携带完整结果与结束原因。
Pi 的事件协议正是这个形状的展开版,本章末尾会看到它的源码。先用一次真实运行把序列看清楚——下面是本章实验第 3 部分的真实输出,一句短回复 你好,我是 Pi 被切成三块:
#1 start
#2 text_start
#3 text_delta delta="你好,"
#4 text_delta delta="我是 "
#5 text_delta delta="Pi"
#6 text_end content="你好,我是 Pi"
#7 done stopReason=stop
序列形状:start → text_start → text_delta ×3 → text_end → done注意这里有两层嵌套的「开始/结束」:外层的 start / done 圈住整轮回复,内层的 text_start / text_end 圈住一个内容块。一轮回复里可以有多个内容块(先一段文字、再一个工具调用),于是内层这对标记会出现多次,外层只有一次。
图 3.3-1 一次流式回复的事件序列
阅读顺序:从上到下,时间向下流动。虚线是模型服务通过网络推来的消息,实线是 Pi 转换后交给上层的统一事件。请关注三点:一、`start` 到达时内容还是空的,它只表示「开始了」;二、中间的 `text_delta` 有任意多个,界面收到一个就打印一个,这就是打字机效果的来源;三、真正带「完整消息」的只有最后的 `done`——中间任何一刻拿到的都是半成品。这三层参与者对应真实源码:模型服务是各家 API,适配层是 `packages/ai/src/api/` 下的文件,界面部分见 [5.4](/pi-runtime/streaming-events)。
数据怎么到你机器上:SSE 一句话
事件在网络上是怎么传的?绝大多数模型 API 用的是 SSE(Server-Sent Events,服务器发送事件):一句话概括——客户端发一个普通的 HTTP 请求,服务器不一次性写完响应体,而是沿着这条一直不关的连接,源源不断地推送一条条文本消息,直到主动结束。
每条消息的格式简单到有点朴素:几行 字段: 值,用一个空行分隔下一条。形状大致如下(这是 SSE 的通用格式示意,具体字段名各家 API 自己定):
event: content_block_delta
data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"你好"}}
event: content_block_delta
data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":",我是"}}Pi 里这条链路分成几段清清楚楚的转换,每一段都是 2.6 学过的异步生成器:
图 3.3-2 从网络字节到统一事件:Pi 的流式数据链路
阅读顺序:从左到右,每个方框是一层转换,数据一块进、一块出,全程没有任何一处「攒齐了再处理」。前三层都在 `packages/ai/src/api/anthropic-messages.ts` 里:`iterateSseMessages`(第 387 行)把字节流切成 SSE 消息,`iterateAnthropicEvents`(第 446 行)把消息解析成该家 API 的原始事件,第 573 行起的循环把原始事件翻译成 Pi 自己的统一事件。第四层的 `EventStream` 是 [2.6](/foundations/async-iterators) 见过的「推转拉」缓冲。最右边的消费方式是一句 `for await`——本章下一节就写它。
这一层层转换的意义在于:上层完全不知道下面是哪家 API。换一家模型服务,只是换掉前三层里的解析代码,第四层往右的所有代码——Agent 循环、界面、Session 存储——一行都不用改。这是 6.2 会展开的 Provider(模型服务提供方)抽象。
兑现 2.6:流式就是异步迭代器
现在把 2.6 那把钥匙插进锁孔。回忆那张四格表:「多个值 × 异步到来」这一格填的是 AsyncIterable<T>。而流式回复正是:多个事件,一个一个异步到来,最后结束。两者是同一个形状——所以 Pi 里表达一次流式请求的返回值,就是一个可以被 for await..of 消费的对象。
生产端用 async function*。下面是本章实验里的假模型(FakeModel,不联网、不需要任何 API Key 的冒牌模型)——它是上一章那个「一次性返回完整回复」的假模型的流式升级版,把一段写死的文本切块,隔一会儿吐一块:
export async function* streamText(text: string, options: StreamOptions = {}): AsyncGenerator<StreamEvent> {
const { delayMs = 60, chunkSize = 4 } = options;
yield { type: "start" };
yield { type: "text_start" };
let accumulated = "";
for (const chunk of splitIntoChunks(text, chunkSize)) {
await sleep(delayMs); // 模拟模型生成这一块 + 网络传输的时间
accumulated += chunk;
yield { type: "text_delta", delta: chunk };
}
yield { type: "text_end", content: accumulated };
yield { type: "done", message: { role: "assistant", text: accumulated, stopReason: "stop" } };
}消费端就是一句 for await,循环体里做的事简单到不像话——收到一块,打印一块:
for await (const event of streamText(REPLY, { delayMs: 60, chunkSize: 4 })) {
if (event.type === "text_delta") {
process.stdout.write(event.delta); // write 不自动换行,块接着块,形成打字机效果
}
}实验里把「非流式」和「流式」两种方式各跑一遍,真实输出是:
== 第 1 部分:非流式——整段到齐才看得见 ==
流式输出的意思是:模型每生成一小块,就立刻把这一小块发出来。
看到第一个字用了 494ms,整段收完也是 494ms(同一时刻)
== 第 2 部分:流式——第一块几十毫秒就到 ==
流式输出的意思是:模型每生成一小块,就立刻把这一小块发出来。
看到第一个字只用了 62ms,整段收完 498ms两行数字把本章第一节的道理量化了:总耗时几乎相同(494ms 对 498ms),首字时间差了近 8 倍。这个比例只取决于回复被切成多少块——真实的几十秒长回复会被切成几百块,差距只会更悬殊。
还有一个容易被忽略的事实:实验里的「非流式」函数 completeText 内部用的就是 streamText,只是把所有块攒齐了才 return。非流式接口可以由流式接口实现,反过来不行。所以底层协议只提供流式是完全合理的——想要哪种,上层自己决定。
汇总:把 delta 拼回一条完整消息
界面要的是「每来一块就显示一块」,但会话历史里存的必须是一条完整消息——下一轮请求要把它原样发回给模型(3.2 讲过上下文是怎么攒起来的)。所以每个流式消费者都要做汇总:遍历事件,把 text_delta 一块块拼起来。
实验的 collect 函数做的正是这件事,它同时留了个心眼——把自己拼出来的文本和 done 事件里的完整消息比一比:
== 第 4 部分:把事件流汇总成一条完整消息 ==
事件总数:12
其中 text_delta:8 个
自己拼出来的文本:流式输出的意思是:模型每生成一小块,就立刻把这一小块发出来。
与 done 事件里的完整消息一致? true这个 true 是流式协议的一条不变量:把所有 delta 按顺序拼起来,必须等于终局事件里的完整内容。任何 Provider 实现都得保证它,否则界面上看到的和历史里存的就会对不上。
Pi 在这件事上更进一步:它的每个 delta 事件除了 delta 之外,还带一个 partial 字段,装着「到此刻为止的半成品消息」——消费者可以直接渲染 partial,不必自己维护累加变量。从源码结构看,这个 partial 每次指向的都是同一个对象(packages/ai/src/api/anthropic-messages.ts 第 495 行创建的 output,随后被不断修改)。一种看法是:这样最省内存也最省事,界面随时能拿到最新全貌;代价是消费者绝不能把某个事件里的 partial 存下来当「那一刻的快照」——它会跟着后续事件一起变。
「部分 JSON」问题:流式的第一个真麻烦
现在把镜头转向工具调用。模型要求调用一个工具时,得给出参数,参数在网络上是一段 JSON 文本,比如 {"city":"北京","days":3}。既然是流式,这段文本当然也是一块一块到的。麻烦立刻出现——半截 JSON 不是合法 JSON。实验第 5 部分把这个过程一帧一帧打了出来(真实输出):
工具名先到:get_weather(参数一个字都还没来)
第 1 块后:{"city → JSON.parse 失败:文本还没写完
第 2 块后:{"city":"北京" → JSON.parse 失败:文本还没写完
第 3 块后:{"city":"北京","days → JSON.parse 失败:文本还没写完
第 4 块后:{"city":"北京","days":3} → JSON.parse 成功最朴素的做法是「等 toolcall_end 到了再解析」——正确,但浪费了流式的好处:参数长一点,用户在收完前什么都看不到。Pi 的做法是每收到一块就尽力解析一次:用 partial-json 这个第三方库把残缺文本解析成「半成品对象」,{"city":"北京","days 会被解析成 { city: "北京" },界面于是能立刻显示「正在查询北京的天气」。
Pi 的封装在 packages/ai/src/utils/json-parse.ts 的 parseStreamingJson(第 104 行):先试标准 JSON.parse,失败了再试修复转义字符,再失败才交给 partial-json,全都失败就返回空对象——永远不抛异常,因为流式期间「解析不出来」是常态而不是错误。完整的工具调用流程、参数校验与执行留到 3.4 和 6.1、6.4。
Pi 中哪里用到了它
第一处:事件协议本身。Pi 把「一次流式回复由哪些事件组成」写成了一个可辨识联合,这是全书后面反复出现的核心类型之一:
AssistantMessageEvent/**
* Event protocol for AssistantMessageEventStream.
*
* Streams should emit `start` before partial updates, then terminate with either:
* - `done` carrying the final successful AssistantMessage, or
* - `error` carrying the final AssistantMessage with stopReason "error" or "aborted"
* and errorMessage.
*/
export type AssistantMessageEvent =
| { type: "start"; partial: AssistantMessage }
| { type: "text_start"; contentIndex: number; partial: AssistantMessage }
| { type: "text_delta"; contentIndex: number; delta: string; partial: AssistantMessage }
| { type: "text_end"; contentIndex: number; content: string; partial: AssistantMessage }
// …(省略:thinking_start / thinking_delta / thinking_end 三个成员,形状与 text 三件套一致)
| { type: "toolcall_start"; contentIndex: number; partial: AssistantMessage }
| { type: "toolcall_delta"; contentIndex: number; delta: string; partial: AssistantMessage }
| { type: "toolcall_end"; contentIndex: number; toolCall: ToolCall; partial: AssistantMessage }
| { type: "done"; reason: Extract<StopReason, "stop" | "length" | "toolUse">; message: AssistantMessage }
| { type: "error"; reason: Extract<StopReason, "aborted" | "error">; error: AssistantMessage };12 个成员,正好是本章讲的形状:一对外层边界(start / done),三组内层三件套(文本、思考、工具调用),外加一个表达「中途失败」的 error。几个细节值得指出:
- 源码里的协议注释是官方说明:流必须先发
start,并且一定以done或error之一收尾——消费者可以依赖这一点写循环,不必担心流悄无声息地断掉。 contentIndex回答了「这一块属于第几个内容块」。一轮回复里有多个内容块时,靠它把 delta 归位——本章开头说「纯字符串流表达不了归属」,这就是补上的那一维。thinking_*三件套对应部分模型的「思考」内容,形状与文本完全一致,所以本章不单独讲。done与error的reason用了Extract<...>:从StopReason这个 union 里挑出允许的几个值,把「成功时不可能是 aborted」这种规则直接写进了类型(2.3 的思路,2.4 的工具类型)。
第二处:事件是在哪里产生的。翻译原始 SSE 事件的循环里,工具调用参数那一段同时展示了「产生 delta 事件」和「部分 JSON」两件事:
parseStreamingJson} else if (event.delta.type === "input_json_delta") {
const index = blocks.findIndex((b) => b.index === event.index);
const block = blocks[index];
if (block && block.type === "toolCall") {
block.partialJson += event.delta.partial_json;
block.arguments = parseStreamingJson(block.partialJson);
stream.push({
type: "toolcall_delta",
contentIndex: index,
delta: event.delta.partial_json,
partial: output,
});
}
}三行代码三件事:累加原始文本(partialJson)、尽力解析成对象(arguments)、把统一事件推给上层。stream.push 里的 stream 就是 2.6 讲过的 EventStream——网络这头推进去,Agent 那头 for await 拉出来。整条链路怎么继续走到终端界面上,是 5.4 流式事件如何传播到界面。
实践任务
labs/agent-concepts/02-streaming目标:亲手做一遍本章全部内容——非流式与流式的体验对比、用 async function* 逐块吐事件、for await 打字机、事件序列观察、汇总回完整消息、部分 JSON 现场。实验目录:labs/agent-concepts/02-streaming(全书实验索引见实践任务索引)。全程不联网、不需要任何 API Key。
步骤:
进入实验目录,安装依赖并运行:
shcd labs/agent-concepts/02-streaming npm install npm start盯着第 2 部分的输出:文字是一块一块长出来的。对照
src/fake-model.ts里的streamText,指出「等一会儿」和「交出一个事件」分别是哪一行。对比第 1、2 部分的两组毫秒数,用一句话说出流式改善了什么、没有改善什么。
把
main.ts里第 2 部分的chunkSize从 4 改成 1,再改成 20,重新运行:观察打字机的颗粒感和首字时间怎么变,解释为什么首字时间随chunkSize变大而变长。在第 3 部分的循环里,把
streamText换成streamToolCall("get_weather", '{"city":"上海"}', { delayMs: 0, chunkSize: 4 }),观察序列形状从文本三件套变成工具调用三件套——这正是 Pi 的AssistantMessageEvent里两组事件的关系。思考题:第 5 部分为什么前三行必然
JSON.parse失败?如果换成 Pi 用的partial-json,第 2 块之后能解析出什么?
预期现象:npm start 的输出与实验目录下 expected-output.txt 一致(第 1、2 部分的四处毫秒数每次运行略有浮动,属正常现象),关键几行:
== 第 1 部分:非流式——整段到齐才看得见 ==
看到第一个字用了 494ms,整段收完也是 494ms(同一时刻)
== 第 2 部分:流式——第一块几十毫秒就到 ==
看到第一个字只用了 62ms,整段收完 498ms
序列形状:start → text_start → text_delta ×3 → text_end → done
与 done 事件里的完整消息一致? true
第 3 块后:{"city":"北京","days → JSON.parse 失败:文本还没写完如何判断成功:肉眼看到第 2 部分的打字机效果;能不看代码说出事件序列形状;能指出哪个事件携带完整消息;第 4 步能解释首字时间的变化。
常见错误:
- 整行一次性出现看不到打字机:多半把输出重定向到了文件或管道(
npm start > out.txt、npm start | cat),换成终端里直接运行; - 把
for await写成for:for..of走同步迭代协议,编译器会提示对象上没有Symbol.iterator; - 在
switch (event.type)里漏掉default分支:本实验里不会报错,但真实代码里漏处理一类事件往往表现为「界面偶尔少显示一段内容」,很难查。
对应源码位置:packages/ai/src/types.ts 第 493–513 行的 AssistantMessageEvent(实验里的 StreamEvent 就是它的缩小版);packages/ai/src/api/anthropic-messages.ts 第 654–666 行(实验第 5 部分的真实对应物)。
本章小结
- 流式输出不加快模型生成,只把「第一个词到达的时间」从整段耗时缩短到一小块的耗时;它换来的是可感知的响应速度、中途取消的可能,以及边到边处理的能力。
- 流式结果的形状不是一串字符,而是一串带类型的事件:外层
start/done圈住整轮回复,内层start/delta/end圈住一个内容块;完整消息与结束原因只在终局事件里。 - SSE(Server-Sent Events)是模型 API 传输这串事件的常用方式:一个普通 HTTP 请求,服务器沿着不关闭的连接推送一条条文本消息。Pi 把它解码、解析、翻译成自己的统一事件,中间每一层都是异步生成器。
- 流式与 2.6 的异步迭代器是同一个形状:生产端
async function*,消费端for await..of。非流式接口可以由流式实现,反之不行。 - 消费者有两个职责:实时显示每个 delta、把所有 delta 汇总回一条完整消息;「拼出来的内容 == 终局事件里的内容」是协议的不变量。Pi 还在每个事件上带一个
partial半成品消息,但它是同一个被反复修改的对象,不能当快照保存。 - 工具调用参数是一段流式到达的 JSON 文本,中途不是合法 JSON。Pi 用
partial-json每收到一块就尽力解析出半成品对象,解析失败返回空对象而不抛异常。
关键术语:流式输出(Streaming)、事件(Event)、增量(delta)、终局事件(done / error)、SSE(Server-Sent Events,服务器发送事件)、可辨识联合(Discriminated Union)、异步迭代器(Async Iterator)、部分 JSON(Partial JSON)、Provider(模型服务提供方)
关键源码索引:packages/ai/src/types.ts 第 493–513 行 AssistantMessageEvent;packages/ai/src/api/anthropic-messages.ts 的 iterateSseMessages(387 行)、iterateAnthropicEvents(446 行)、事件翻译循环(573 行起)与工具参数分支(654–666 行);packages/ai/src/utils/json-parse.ts 的 parseStreamingJson(104 行);packages/ai/src/utils/event-stream.ts 的 EventStream
自测问题:
- 有人说「开了流式,模型回复就更快了」。这句话哪里对、哪里不对?请分别用「总耗时」和「首字时间」两个指标回答。
- 为什么流式协议要发
text_start/text_end这样的边界事件,只发text_delta不行吗?举一个只发 delta 就会出错的场景。 - 一次回复里,哪个事件携带「权威的完整消息」?如果消费者只保存了每个
delta,它还需要做什么才能把这轮回复存进会话历史? - 工具调用参数
{"path":"/etc/hosts","limit":100}流式到达,某一刻收到的文本是{"path":"/etc/ho。直接JSON.parse会怎样?Pi 会怎么处理,界面能显示出什么?
下一章预告:3.4 工具调用(Tool Calling)——本章末尾那个 toolcall_delta 只是开头。模型凭什么知道有哪些工具可用?它「要求调用工具」到底是什么形式?参数收完之后谁去执行、结果又怎么送回模型?下一章把这条链路走完,Agent 也就从「会聊天」变成了「会干活」。