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 表示流是否结束,value 是 Uint8Array,需要用 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 通常会逐层转换:
不同接入方式经过的处理层次不同:
- 客户端调用模型 API 通常直接消费供应商的协议事件或响应字段,可以不建立统一的业务事件层;
- 服务端中转会把供应商事件映射成应用统一的业务事件;
- SDK 或 LangChain 已经返回结构化增量时,可以跳过协议解析,直接映射成业务事件。
封装 SSE 流处理方法
服务端处理 SSE 时会反复涉及字节解码、事件组装和响应生命周期。将这些逻辑封装后,业务代码只需读取上游事件、转换字段并写出业务事件。这里仅解析本文使用的 event 字段和 JSON 格式的 data 字段,不覆盖完整的 SSE 规范。
后面的示例都建立在这三个函数上:
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()
}
}
lib/sse.ts
const encoder = new TextEncoder()
export function writeSSE(
controller: ReadableStreamDefaultController<Uint8Array>,
event: string,
data: unknown
) {
controller.enqueue(
encoder.encode(`event: ${event}\ndata: ${JSON.stringify(data)}\n\n`)
)
}
lib/sse.ts
export function createSSEResponse(
request: Request,
produce: (
controller: ReadableStreamDefaultController<Uint8Array>,
signal: AbortSignal
) => Promise<void>
) {
const cancelController = new AbortController()
const signal = AbortSignal.any([request.signal, cancelController.signal])
let cancelled = false
const stream = new ReadableStream({
async start(controller) {
try {
await produce(controller, signal)
if (!signal.aborted) {
writeSSE(controller, 'done', {})
}
} catch (error) {
if (!signal.aborted) {
const message = error instanceof Error ? error.message : '执行失败'
writeSSE(controller, 'error', { message })
}
} finally {
if (!cancelled) {
controller.close()
}
}
},
cancel() {
cancelled = true
cancelController.abort()
},
})
return new Response(stream, {
headers: {
'Content-Type': 'text/event-stream; charset=utf-8',
'Cache-Control': 'no-cache, no-transform',
},
})
}
模型调用方式
模型流可以由客户端直接消费,也可以先在服务端完成转换,区别在于密钥、鉴权和供应商协议由哪一端管理。
客户端调用模型 API
客户端直接调用模型 API 时,可以使用 parseSSE 解析响应体,再提取文本增量并更新界面。下面的示例直接处理 Chat Completions 的 SSE 数据,因此客户端需要了解上游响应格式。这种方式适合 demo 和由用户自行管理 API Key 的本地 Agent,示例中的 apiKey 和 model 由用户在本地配置。
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、error 和 done 三类事件,客户端可以先用一个事件处理器统一分发:
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;自定义事件(如 done、tool_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 请求可以使用 AbortController,EventSource 则调用 close()。服务端也要感知连接关闭并中止上游调用,避免模型继续生成或向已关闭的连接写入数据。
createSSEResponse 会合并请求断开和响应流取消产生的信号,两者任一发生都会中止上游调用。业务代码只需将 signal 透传给上游 fetch、SDK 调用或 model.stream()。
将错误作为事件返回
流式接口的错误可能发生在响应开始前或传输过程中。请求校验、鉴权失败等已知错误应在流开始前通过 HTTP 状态码返回;响应开始后,状态码已经发出,上游模型断开或工具超时等错误只能转换成客户端可识别的事件。
event: error
data: {"message":"模型连接中断"}
客户端收到 error 事件后,应停止加载、展示原因并允许用户重试。createSSEResponse 会在 produce 抛错时写出该事件,业务代码只需保留具体的错误上下文。
总结
流式处理的核心是将网络片段逐层解析为协议事件和业务事件,再由客户端按事件更新 UI。服务端可以根据密钥、鉴权和供应商适配需求选择不同的模型调用方式,同时处理取消、错误、心跳和代理缓冲,确保事件及时到达并正确结束连接。