Agent 应用中的流处理

Agent 请求通常不是一次模型调用,而是由模型、检索和工具执行组成的连续过程。服务端需要把这些步骤产生的增量结果转换为稳定的业务事件,再交给客户端逐步更新回答和执行状态。本文以 SSE 为例,介绍 Agent 应用中的流读取、服务端中转和客户端消费。

什么是流

流(stream)描述数据按产生顺序逐步到达的过程。生产方产生一部分数据后即可发送,消费方收到后立即处理,不必等待完整结果。普通响应则要等结果全部生成后,再一次性交给消费方。

流描述数据如何到达,不限定内容或编码。SSE 定义网络传输中的事件格式,Web Streams 和 Node.js Streams 提供在代码中读取、写入和转换这些数据的 API。

Node.js Streams

Node.js Streams 是 Node.js 运行时提供的流 API,主要包含以下类型:

  • Readable 用于读取数据;
  • Writable 用于写入数据;
  • Duplex 支持双向读写;
  • Transform 用于在读写过程中转换数据。

Node.js Streams 主要用于文件、网络连接和进程输入输出。可读流通常通过 pipe() 连接到其它流,也可以通过 for await...of 按块消费。它与 Web Streams 的职责相近,但类型和方法不同。两者需要互操作时,Node.js 提供 Readable.fromWeb()Readable.toWeb()

Web Streams

Web Streams API 是 Web 平台提供的流 API,主要包含以下类型:

  • ReadableStream 用于逐段读取数据;
  • WritableStream 用于逐段写入数据;
  • TransformStream 用于在读取和写入之间转换数据。

客户端通过 fetch 发起请求后,response.body 是一个 ReadableStream<Uint8Array>。读取时先通过 getReader() 获取读取器,再调用 reader.read() 取得数据片段。done 表示流是否结束,valueUint8Array,需要用 TextDecoder 解码成文本。下面的代码只演示读取一次,完整消费需要在循环中重复调用 reader.read()

const response = await fetch('/api/chat', {
  method: 'POST',
  headers: { 'Content-Type': 'application/json' },
  body: JSON.stringify({ messages }),
})

if (!response.body) {
  throw new Error('模型响应体为空')
}

const reader = response.body.getReader()
const decoder = new TextDecoder()
const { value, done } = await reader.read()
const text = done ? decoder.decode() : decoder.decode(value, { stream: true })

一次读取只能取得当前可用的字节,不保证对应完整消息。连续读取时要复用同一个 TextDecoder,中间片段使用 { stream: true } 保留未完成的字节,流结束时再调用 decode() 清空内部缓冲。

流数据的处理层次

流数据从网络到 UI 通常会逐层转换:

层次说明
网络片段reader.read() 返回的 value 是本次读取到的字节,不保证对应一条完整消息。
文本片段TextDecoder 解码后的文本,仍可能只包含半条消息或多条消息的一部分。
协议事件(可选)按 SSE 空行边界拼出的完整事件,event 名称由发送方约定。
业务事件将上游结果映射为应用统一的事件类型,例如文本增量、工具状态、错误和结束信号。
UI 状态业务事件处理器根据事件类型更新界面,例如追加回答、更新状态、展示错误或完成收尾。

不同接入方式经过的处理层次不同:

  • 客户端调用模型 API 通常直接消费供应商的协议事件或响应字段,可以不建立统一的业务事件层;
  • 服务端中转会把供应商事件映射成应用统一的业务事件;
  • SDK 或 LangChain 已经返回结构化增量时,可以跳过协议解析,直接映射成业务事件。

封装 SSE 流处理方法

服务端处理 SSE 时会反复涉及字节解码、事件组装和响应生命周期。将这些逻辑封装后,业务代码只需读取上游事件、转换字段并写出业务事件。这里仅解析本文使用的 event 字段和 JSON 格式的 data 字段,不覆盖完整的 SSE 规范。

后面的示例都建立在这三个函数上:

函数职责
parseSSEReadableStream 解析成 SSE 事件流,调用方用 for await...of 消费。
writeSSE服务端向 controller 写入一条 SSE 事件。
createSSEResponse包装一个标准的 SSE Response,统一处理结束、错误和取消信号。
parseSSE
writeSSE
createSSEResponse
lib/sse.ts
export async function* parseSSE(body: ReadableStream<Uint8Array>) {
  const reader = body.getReader()
  const decoder = new TextDecoder()
  let buffer = ''
  let finished = false

  try {
    while (true) {
      const { value, done } = await reader.read()

      if (done) {
        buffer += decoder.decode()
        finished = true
        break
      }

      buffer += decoder.decode(value, { stream: true })
      const blocks = buffer.split(/\r?\n\r?\n/)
      buffer = blocks.pop() ?? ''

      for (const block of blocks) {
        const lines = block.split(/\r?\n/)
        const event =
          lines
            .find((line) => line.startsWith('event:'))
            ?.slice('event:'.length)
            .trim() || 'message'
        const data = lines
          .filter((line) => line.startsWith('data:'))
          .map((line) => line.slice('data:'.length).replace(/^ /, ''))
          .join('\n')

        if (data) {
          yield { event, data }
        }
      }
    }
  } finally {
    if (!finished) {
      try {
        await reader.cancel()
      } catch {
        // 流已经出错或被取消时,cancel 可能再次失败。
      }
    }
    reader.releaseLock()
  }
}

模型调用方式

模型流可以由客户端直接消费,也可以先在服务端完成转换,区别在于密钥、鉴权和供应商协议由哪一端管理。

客户端调用模型 API

客户端直接调用模型 API 时,可以使用 parseSSE 解析响应体,再提取文本增量并更新界面。下面的示例直接处理 Chat Completions 的 SSE 数据,因此客户端需要了解上游响应格式。这种方式适合 demo 和由用户自行管理 API Key 的本地 Agent,示例中的 apiKeymodel 由用户在本地配置。

import { parseSSE } from '@/lib/sse'

const response = await fetch('https://api.openai.com/v1/chat/completions', {
  method: 'POST',
  headers: {
    Authorization: `Bearer ${apiKey}`,
    'Content-Type': 'application/json',
  },
  body: JSON.stringify({ model, messages, stream: true }),
})

if (!response.ok) {
  throw new Error(`模型请求失败:${response.status}`)
}

if (!response.body) {
  throw new Error('模型响应体为空')
}

for await (const { data } of parseSSE(response.body)) {
  if (data === '[DONE]') break

  const payload = JSON.parse(data)
  const content = payload.choices?.[0]?.delta?.content
  if (content) appendAssistantText(content)
}

服务端中转模型 API

本文使用 Next.js Route Handler 编写服务端示例,但流处理本身不依赖 Next.js,同样适用于其它 Node.js 服务。

服务端中转适合需要集中管理密钥、鉴权和请求控制的场景。它接收客户端请求,调用模型 API,再将 Chat Completions 的文本增量映射为统一的 message 事件。客户端只需处理这个业务事件,不必依赖上游响应字段。

app/api/chat/route.ts
import { createSSEResponse, parseSSE, writeSSE } from '@/lib/sse'

export async function POST(request: Request) {
  const { messages } = await request.json()

  return createSSEResponse(request, async (controller, signal) => {
    const response = await fetch('https://api.openai.com/v1/chat/completions', {
      method: 'POST',
      headers: {
        Authorization: `Bearer ${process.env.OPENAI_API_KEY}`,
        'Content-Type': 'application/json',
      },
      body: JSON.stringify({
        model: process.env.OPENAI_MODEL!,
        messages,
        stream: true,
      }),
      signal,
    })

    if (!response.ok) {
      throw new Error(`模型请求失败:${response.status}`)
    }

    for await (const { data } of parseSSE(response.body!)) {
      if (data === '[DONE]') break

      const payload = JSON.parse(data)
      const content = payload.choices?.[0]?.delta?.content
      if (content) writeSSE(controller, 'message', { content })
    }
  })
}

服务端通过 SDK 调用模型

使用供应商 SDK 时,请求体、鉴权和流式响应解析由 SDK 处理。服务端用 for await...of 消费返回的流,再将需要的 Responses 事件映射为业务事件。与前一个示例相比,这里不再调用 parseSSE,因为 SDK 已经完成了协议解析。

app/api/chat/route.ts
import OpenAI from 'openai'
import { createSSEResponse, writeSSE } from '@/lib/sse'

const client = new OpenAI({
  apiKey: process.env.OPENAI_API_KEY,
})

export async function POST(request: Request) {
  const { messages } = await request.json()

  return createSSEResponse(request, async (controller, signal) => {
    const stream = await client.responses.create(
      {
        model: process.env.OPENAI_MODEL!,
        input: messages,
        stream: true,
      },
      { signal }
    )

    for await (const event of stream) {
      if (event.type === 'response.output_text.delta') {
        writeSSE(controller, 'message', { content: event.delta })
      }

      if (event.type === 'response.failed') {
        throw new Error(event.response.error?.message ?? '模型生成失败')
      }

      if (event.type === 'error') {
        throw new Error(event.message)
      }
    }
  })
}

服务端通过 LangChain 调用模型

如果希望统一不同供应商的调用方式,或需要接入工具和多步 Agent,可以使用 LangChain。它同样返回可异步迭代的增量,服务端消费 model.stream() 的结果,并将其映射为 message 等业务事件。这样切换供应商时,通常只需替换对应的模型类和配置,调用方式和业务事件映射可以保持不变。

app/api/chat/route.ts
import { ChatOpenAI } from '@langchain/openai'
import { createSSEResponse, writeSSE } from '@/lib/sse'

const model = new ChatOpenAI({
  apiKey: process.env.OPENAI_API_KEY,
  model: process.env.OPENAI_MODEL!,
})

export async function POST(request: Request) {
  const { messages } = await request.json()

  return createSSEResponse(request, async (controller, signal) => {
    const modelStream = await model.stream(messages, { signal })

    for await (const chunk of modelStream) {
      if (chunk.text) {
        writeSSE(controller, 'message', { content: chunk.text })
      }
    }
  })
}

客户端消费业务事件流

采用服务端中转时,原始 API、SDK 或 LangChain 返回的结果都会先转换为业务事件,客户端再按事件类型分别处理。统一事件可以覆盖文本、工具和生命周期状态:

事件类型含义客户端行为
message模型文本增量追加到回答区域
tool_call工具调用开始或更新展示工具调用状态
tool_result工具执行结果展示结果或折叠详情
statusAgent 当前执行状态更新状态文案
error生成或执行失败展示错误并结束加载状态
done流结束收尾并允许继续输入

不同事件不能全部作为文本追加。前面的文本流示例只写出了 messageerrordone 三类事件,客户端可以先用一个事件处理器统一分发:

function handleAgentEvent(event) {
  switch (event.type) {
    case 'message':
      appendAssistantText(event.content)
      break

    case 'error':
      showError(event.message)
      finishStreaming()
      break

    case 'done':
      finishStreaming()
      break
  }
}

用 EventSource 消费 SSE

EventSource 通过 GET 建立服务端到浏览器的单向 SSE 连接,并按 event: 名称自动分发事件。它适合不需要 POST body 或自定义请求头的场景。下面的示例假设服务端已经提供了对应的 GET SSE 接口,并用 sessionId 标识要消费的会话;前面的 POST Route Handler 不能直接作为这个接口使用。

const source = new EventSource(
  `/api/chat?sessionId=${encodeURIComponent(sessionId)}`
)

source.addEventListener('message', (e) => {
  handleAgentEvent({ type: 'message', ...JSON.parse(e.data) })
})

source.addEventListener('done', () => {
  handleAgentEvent({ type: 'done' })
  source.close()
})

source.addEventListener('error', (event) => {
  if (event instanceof MessageEvent && event.data) {
    handleAgentEvent({ type: 'error', ...JSON.parse(event.data) })
    source.close()
    return
  }

  // 连接异常不调用 close(),EventSource 会自动重连。
})

未指定 event: 时,事件名默认为 message;自定义事件(如 donetool_call)需要显式监听。业务错误和连接异常都会触发 error 监听器,但业务错误事件带有 data,连接异常通常没有。收到业务错误后调用 source.close() 结束连接,连接异常则交给 EventSource 自动重连。

用 fetch 消费 SSE

EventSource 只支持 GET,且不能直接携带自定义请求头。Chat 或 Agent 接口通常需要通过 POST body 传递消息,还可能携带鉴权信息和请求上下文。这类接口可以用 fetch 发起请求,再通过 parseSSE 解析响应流。

import { parseSSE } from '@/lib/sse'

const response = await fetch('/api/chat', {
  method: 'POST',
  headers: { 'Content-Type': 'application/json' },
  body: JSON.stringify({ messages }),
})

if (!response.ok) {
  handleAgentEvent({ type: 'error', message: `请求失败:${response.status}` })
  return
}

if (!response.body) {
  handleAgentEvent({ type: 'error', message: '响应体为空' })
  return
}

for await (const { event, data } of parseSSE(response.body)) {
  const payload = JSON.parse(data)
  handleAgentEvent({ ...payload, type: event })
}

用 Axios 消费 SSE

如果项目已经使用 Axios,浏览器端可以通过 XHR 的 onDownloadProgress 读取不断累积的 responseText,再自行拆分 SSE 事件,并暂存未完整接收的片段。

import axios from 'axios'

let readOffset = 0
let buffer = ''

function handleSSEBlock(block) {
  const lines = block.split(/\r?\n/)
  const event =
    lines
      .find((line) => line.startsWith('event:'))
      ?.slice('event:'.length)
      .trim() || 'message'
  const data = lines
    .filter((line) => line.startsWith('data:'))
    .map((line) => line.slice('data:'.length).replace(/^ /, ''))
    .join('\n')

  if (!data) return

  const payload = JSON.parse(data)
  handleAgentEvent({ ...payload, type: event })
}

function consumeResponseText(responseText) {
  const chunk = responseText.slice(readOffset)
  readOffset = responseText.length

  buffer += chunk
  const blocks = buffer.split(/\r?\n\r?\n/)
  buffer = blocks.pop() ?? ''

  for (const block of blocks) {
    if (block.trim()) handleSSEBlock(block)
  }
}

try {
  await axios.post(
    '/api/chat',
    { messages },
    {
      responseType: 'text',
      onDownloadProgress({ event }) {
        consumeResponseText(event.target.responseText)
      },
    }
  )
} catch {
  handleAgentEvent({ type: 'error', message: '请求失败' })
}
使用建议

Axios 的下载进度回调可能被节流,文本增量可能成批到达。如果不依赖 Axios 的实例配置或拦截器,建议使用 fetch + parseSSE

流处理的工程实践

上述方法封装了字节、协议和生命周期处理,但事件能否及时到达客户端还取决于部署链路。反向代理、负载均衡器或托管平台可能缓冲响应,导致服务端持续写入,客户端却迟迟收不到事件。部署时应确认整条链路支持流式传输,并关闭代理缓冲。使用 Nginx 时,服务端可以返回 X-Accel-Buffering: no 响应头,关闭当前响应的代理缓冲。

工具执行期间如果长时间没有业务事件,还应发送心跳,并检查各层的空闲超时设置。SSE 注释行不会触发业务事件,可用作心跳,例如 : ping\n\n

处理用户取消和连接关闭

用户停止生成或离开页面时,客户端应取消当前连接。fetch 和 Axios 请求可以使用 AbortControllerEventSource 则调用 close()。服务端也要感知连接关闭并中止上游调用,避免模型继续生成或向已关闭的连接写入数据。

createSSEResponse 会合并请求断开和响应流取消产生的信号,两者任一发生都会中止上游调用。业务代码只需将 signal 透传给上游 fetch、SDK 调用或 model.stream()

将错误作为事件返回

流式接口的错误可能发生在响应开始前或传输过程中。请求校验、鉴权失败等已知错误应在流开始前通过 HTTP 状态码返回;响应开始后,状态码已经发出,上游模型断开或工具超时等错误只能转换成客户端可识别的事件。

event: error
data: {"message":"模型连接中断"}

客户端收到 error 事件后,应停止加载、展示原因并允许用户重试。createSSEResponse 会在 produce 抛错时写出该事件,业务代码只需保留具体的错误上下文。

总结

流式处理的核心是将网络片段逐层解析为协议事件和业务事件,再由客户端按事件更新 UI。服务端可以根据密钥、鉴权和供应商适配需求选择不同的模型调用方式,同时处理取消、错误、心跳和代理缓冲,确保事件及时到达并正确结束连接。