Skip to content

Node.js 接入大模型 SSE 流式输出全解析 ​

为什么大模型需要流式输出? ​

大语言模型(LLM)的推理过程是自回归生成的——模型一个 Token 一个 Token 逐一预测输出。 如果采用传统 HTTP 同步等待方式,一次 1000 Token 的长回复可能需要等待 10~20 秒,客户端在此期间处于完全白屏无响应状态,首字延迟(TTFT, Time To First Token)极差。

采用 Server-Sent Events (SSE) 流式传输,模型生成首个 Token(几百毫秒内)即可立即推送到前端进行打字机渲染,极大提升用户体验。


SSE 协议数据格式详解 ​

SSE 是基于标准 HTTP/1.1 或 HTTP/2 的长连接单向推送协议。服务端响应头关键字段为:

http
Content-Type: text/event-stream; charset=utf-8
Cache-Control: no-cache
Connection: keep-alive

每一条事件消息以两组连续换行符 \n\n 分隔。典型的 LLM 输出流如下:

text
data: {"id":"1","choices":[{"delta":{"content":"你"}}]}

data: {"id":"2","choices":[{"delta":{"content":"好"}}]}

data: [DONE]

核心技术难点与避坑指南 ​

1. 网络数据包 Chunk 与 SSE 消息边界不一致 ​

TCP 和 HTTP 的 chunk 分片取决于网络拥塞控制与 MTU 大小。一次网络 data chunk 并不等同于一条完整的 SSE 消息!

  • 一个 chunk 可能包含半条 JSON;
  • 一个 chunk 可能合并了 3 条 SSE 消息。

错误做法:直接 JSON.parse(chunk.toString()),这必然会在高并发或较长消息下崩溃。

正确解法:维护一个文本缓冲区(Buffer / Accumulator),按 \n\n 扫描切割,剩余未闭合部分留存至下一个 chunk 拼接。

2. UTF-8 多字节字符截断问题 ​

中文字符在 UTF-8 下通常占用 3 个字节。网络分片可能正好落在一个中文字符的第 2 和第 3 个字节之间。 如果直接使用 Buffer.toString(),会导致截断处变为乱码(``)。 解决办法:使用 Node.js 内置的 TextDecoder({ fatal: false, stream: true }),它会在内部自动缓冲不完整的字节序列!


生产级消费实现 (TypeScript) ​

typescript
import { Readable } from 'node:stream'

/**
 * 消费大模型 SSE 流的异步生成器
 */
export async function* parseSSEStream(
  stream: ReadableStream<Uint8Array>
): AsyncGenerator<string, void, unknown> {
  const reader = stream.getReader()
  const decoder = new TextDecoder('utf-8', { stream: true })
  let buffer = ''

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

      // stream: true 保证多字节字符不被切割破坏
      buffer += decoder.decode(value, { stream: true })

      const lines = buffer.split('\n')
      // 最后一行可能未接收完整,保留在 buffer 中
      buffer = lines.pop() || ''

      for (const line of lines) {
        const trimmed = line.trim()
        if (!trimmed || trimmed.startsWith(':')) {
          // 忽略空行或心跳注释
          continue
        }

        if (trimmed.startsWith('data: ')) {
          const payload = trimmed.slice(6).trim()
          if (payload === '[DONE]') {
            return
          }

          try {
            const parsed = JSON.parse(payload)
            const textDelta = parsed.choices?.[0]?.delta?.content
            if (textDelta) {
              yield textDelta
            }
          } catch (e) {
            // 如果某一行不是完整 JSON,视情况记录警告
            console.warn('JSON parse error in SSE line:', payload)
          }
        }
      }
    }
  } finally {
    reader.releaseLock()
  }
}

优雅中断机制 (AbortController) ​

用户在 AI 输出到一半时经常点击“停止生成”按钮。如果服务端不停止向上游大模型拉取流,就会持续扣除 Token 费用。

typescript
const controller = new AbortController()

// 当用户断开前端连接或点击停止时
req.on('close', () => {
  controller.abort()
})

const response = await fetch('https://api.example.com/v1/chat/completions', {
  method: 'POST',
  headers: {
    'Authorization': `Bearer ${process.env.API_KEY}`,
    'Content-Type': 'application/json'
  },
  body: JSON.stringify({
    model: 'gemini-1.5-pro',
    stream: true,
    messages: [{ role: 'user', content: '写一篇长文' }]
  }),
  signal: controller.signal // 传入 signal 允许主动取消
})

总结 ​

  • 牢记网络分包不等于消息边界,必须使用带状态的缓冲区切割。
  • 始终使用 TextDecoder(..., { stream: true }) 避免多字节中文乱码。
  • 生产环境务必绑定 AbortController,在连接断开时切断上游消耗。

基于 VitePress 构建 | 记录真实开发与 AI 协作过程