跳到正文
astorlm
语言: 简体中文
← 地图

第 6 关

流式输出与事件

模型写回复是几个 token 几个 token 地写。流式输出把每一小段一写出来就交给你,而循环发出的事件会告诉应用的其他部分:智能体正在做什么,就在它做的时候。
1/37 手风琴褶数:
  • user
  • assistant
  • tool_result
演出开始。神谕者负责唱,音符顺着音轨滑下来,弹吉他的人把它们弹给观众听。先来一个普通请求:不用流式输出。

EventBus

问题

调用模型、等整个回复回来,在脚本里没问题。可当有人看着屏幕时就不行了:一段长回复可能要十秒、二十秒,这期间屏幕上什么都没有。智能体让情况更糟,因为一次请求在最终回复出现之前,可能要跑好几轮、调用好几个工具。

而且想跟进进度的不只是人。界面要文本,日志要每一步,计费要每一轮的 token。如果这些都得写进循环里,循环最后就得认识你应用里的每一块屏幕、每一个文件、每一个数据库。

解决方案

两块配套使用的东西。流式输出:让服务商以流的形式返回回复,模型每写出一小段文本就立刻发过来。总时间不变,但第一个字会在零点几秒内出现,而不是等到最后。

事件:循环在一条总线上广播每一步,并不在乎谁在听。你的应用让关心的部分去订阅:屏幕听文本,日志听所有事件,计量表听每一轮的结束。加一个监听者不用改动循环。下面是 astorlm 在一次正常运行中发出的事件:

  • turn_start

    一轮开始:循环马上要调用模型。

  • text_delta

    回复的一小段到了。每一轮有很多个。

  • assistant_message

    模型的消息完整了,包括文本和工具调用。

  • tool_execution_start

    一个工具马上要运行,带着它的输入。

  • tool_execution_end

    工具跑完了:它的输出、有没有失败、花了多久。

  • turn_end

    这一轮结束了,带着这次模型调用用掉的 token。

  • session_end

    这次运行结束了:完成、取消或出错。

工具调用从服务商那里也是流式过来的,参数一次来几个字符。循环会先把它们拼完整再运行,所以在总线上,工具调用是整个出现的,在 assistant_message 里。

角色

还是那群熟悉的角色,这次站上了摇滚舞台。

神谕者 模型
主唱,站在升降台上。它的回复变成音符顺着轨道滑下来,一个词一个音符。
轨道 流
每个 text_delta 是一小把音符。不用流式输出,直到最后什么都不会滑下来。
弹吉他的人 你的应用
音符一到音键就弹出来,歌词随之出现在观众面前。
线缆 EventBus
所有事件都从这里跑过。只有订阅了这个事件的监听者才会亮。
监听者 agent.on(…)
歌词屏('text')、一台录音机('event')和一个 token 计量表(turn_end)。
Astor 循环
照常跑循环,一边走一边在总线上广播每一步。
舞台 工具
台口、灯光控制台和烟火箱:crowd_mood、set_lights 和 pyro。

代码

使用 astorlm:服务商总是以流式返回,所以没有什么需要打开的。在调用 run() 之前用 agent.on() 订阅:'text' 接收每一小段文本,'tool-start' 和 'tool-end' 接收工具事件,'event' 接收所有事件。每次调用都会返回一个用来取消订阅的函数。

从零手写:第 2 关的循环,改两处。请求带上 stream: true,把回复当作 server-sent events 来读,再把一小段一小段拼起来;另外有一个监听者列表,通过同一个 emit() 收到每一步。

import { OpenAIProvider, createLocalAgent, tool } from 'astorlm'
import { z } from 'zod'

// The stage's functions, wrapped as tools.
const crowdMood = tool({
  name: 'crowd_mood',
  description: 'Look at the crowd from the front of the stage: how they feel, and what they came to hear.',
  schema: z.object({}),
  execute: async () => stage.readCrowd(), // your code
})

const setLights = tool({
  name: 'set_lights',
  description: 'Set the color of every light on the stage.',
  schema: z.object({ color: z.enum(['amber', 'red', 'blue', 'white']) }),
  execute: async ({ color }) => stage.lights(color), // your code
})

const pyro = tool({
  name: 'pyro',
  description: 'Fire the pyrotechnics on both sides of the stage. Once per show.',
  schema: z.object({}),
  execute: async () => stage.firePyro(), // your code
})

const agent = await createLocalAgent({
  // Any OpenAI-compatible endpoint: OpenAI, Ollama, LM Studio, vLLM, a proxy…
  provider: new OpenAIProvider({
    baseURL: 'http://localhost:11434/v1', // e.g. Ollama's default address
    model: 'your-model', // e.g. 'llama3.1', 'gpt-4o-mini'
    apiKey: 'YOUR_API_KEY', // local servers usually ignore it
  }),
  tools: [crowdMood, setLights, pyro],
  maxTurns: 10,
})

// Streaming is always on: subscribe before run(). Each listener only hears what it asked for.
// The lyrics screen: every piece of text, the moment it arrives.
agent.on('text', (piece) => lyrics.append(piece))

// The tape deck: every event the loop emits, in order.
agent.on('event', (event) => tape.write(event))

// The meter: the tokens each turn reported.
agent.on('event', (event) => {
  if (event.type === 'turn_end' && event.usage) meter.add(event.usage.inputTokens + event.usage.outputTokens)
})

// Each on() returns a function that unsubscribes.
const stopTape = agent.on('tool-start', ({ name }) => console.log('→', name))

const last = await agent.run('Open the show: read the crowd, set the mood, and greet them.')
stopTape()
console.log(last.content) // the same text, whole, once the run is over

注意事项

  • 流式输出不会让模型变快。回复花的时间和以前一样;变的是人看到第一个字的时间。要测首个 token 的时间,而不只是总时间。
  • 监听者只看,不掌舵。监听者的返回值会被忽略,循环也不会等它。想拦下一个工具,或者改变模型读到的内容,你需要的是 hook:下一关。
  • 让监听者跑得快。astorlm 的总线在循环内部一个接一个地调用它们。慢活儿,比如写数据库,应该放进一个队列,监听者只负责往里推。
  • 坏掉的监听者会悄无声息地失败。astorlm 会接住它的错误,让运行继续,这也意味着没人会知道。在你的监听者里自己记日志。
  • 一小段不一定是一个完整的词。一个 text_delta 可能停在一个词中间,或一张 Markdown 表格中间。先显示已有的内容,等更多内容到了再重画。
  • 让人能叫停。人一旦能看着回复一点点到来,就会想打断一个跑偏的回复。把停止按钮接到 agent.abort() 上。