Skip to content

第 5 章 · 事件流订阅与流式输出

本章目标:掌握 agent.subscribe() 与全部事件类型,理解 message_update / text_delta 增量机制,能构建打字机式实时输出。

5.1 为什么要事件流

LLM 生成一段长回复可能需要十几秒。与其让用户盯着空白屏幕,不如把生成过程实时推送出来——这就是事件流的用途。Agent 把回合内发生的一切都建模为事件:

typescript
// 订阅:返回反订阅函数
const unsubscribe = agent.subscribe((event) => {
  console.log(event.type);
});

// 用完记得取消(组件卸载/会话结束时)
unsubscribe();

监听器按注册顺序被 await——agent_end 的异步处理也会计入回合落定,这为「收尾工作」提供了可靠屏障。

5.2 prompt() 的完整事件序列

一次无工具调用的 prompt("Hello") 会依次发出:

text
prompt("Hello")
├─ agent_start                      # 回合开始
├─ turn_start                       # 一个 turn = 一次 LLM 调用 + 工具执行
├─ message_start { userMessage }    # 用户消息入列
├─ message_end   { userMessage }
├─ message_start { assistantMessage }
├─ message_update { partial... }    # ★ 流式增量,可能成百上千次
├─ message_update { partial... }
├─ message_end   { assistantMessage } # 完整回复
├─ turn_end      { message, toolResults: [] }
└─ agent_end      { messages: [...] } # 整个 run 的最终事件

有工具调用时,turn 之间还会插入 tool_execution_start / update / end(第 6、7 章展开)。

5.3 message_update 与 text_delta

注意嵌套结构:message_update 是 Agent 层的事件,其内部携带 assistantMessageEvent(pi-ai 层的原始流事件):

typescript
// 打字机效果的标准写法
agent.subscribe((event) => {
  if (event.type === "message_update") {
    const inner = event.assistantMessageEvent;
    switch (inner.type) {
      case "text_delta":
        // 只输出新增的文本片段
        process.stdout.write(inner.delta);
        break;
      case "thinking_delta":
        // 推理型模型的思考过程增量
        process.stdout.write(`\x1b[2m${inner.delta}\x1b[0m`);
        break;
    }
  }
});

5.4 事件类型全景表

事件说明关键属性
agent_startAgent 开始处理本次 run
agent_endrun 的最终事件,此后不再有循环事件messages: [...]
turn_start新 turn 开始(1 次 LLM 调用 + 工具执行)
turn_endturn 结束message, toolResults
message_start任意消息开始(user/assistant/toolResult)message
message_update仅 assistant 消息,携带增量事件assistantMessageEvent
message_end消息完成message
tool_execution_start工具开始执行toolCallId, toolName, args
tool_execution_update工具流式进度(工具主动上报时)部分结果
tool_execution_end工具执行完毕toolCallId, result

5.5 实战:一个结构化的终端 UI

typescript
// 组合多种事件构建可读的终端输出
agent.subscribe((event) => {
  switch (event.type) {
    case "turn_start":
      process.stdout.write("\n🤖 ");
      break;
    case "tool_execution_start":
      // 工具开始时给出明确提示
      console.log(`\n🔧 调用工具 ${event.toolName} 参数=${JSON.stringify(event.args)}`);
      break;
    case "tool_execution_end":
      console.log(`✅ 工具完成 (${event.toolCallId})`);
      break;
    case "message_update": {
      const inner = event.assistantMessageEvent;
      if (inner.type === "text_delta") process.stdout.write(inner.delta);
      break;
    }
    case "agent_end":
      console.log(`\n\n📊 本回合共 ${event.messages.length} 条消息`);
      break;
  }
});

await agent.prompt("搜索一下 Node.js 22 的新特性并总结");

5.6 本章小结

  • subscribe() 按注册顺序 await 监听器;返回值用于取消订阅;
  • 事件层级:agent → turn → message → tool execution;
  • message_update.assistantMessageEvent.text_delta 是实现打字机效果的钥匙;
  • agent_end 是 run 的屏障事件,适合做持久化等收尾工作。

🧪 随堂测验

点击你认为正确的选项。答错时会展示正确答案与原因解析。

1. 要实现打字机效果,应该监听哪个事件并读取哪个字段?

2. 关于 agent_end 事件,正确的说法是?

3. 工具开始执行时会发出哪个事件?

4. `agent.subscribe()` 的返回值是什么?

🛠️ 动手实践

  1. 给第 4 章的 Agent 加上完整的彩色终端 UI:思考灰色斜体、正文白色、工具黄色提示。
  2. 统计一次含工具调用的回合中各事件类型的出现次数,验证事件序列图。
  3. 注册两个订阅者(一个打印、一个计数),验证它们都按注册顺序收到事件。

下一章第 6 章进入本课程最核心的主题:工具的定义与执行。