前端开发··2 阅读·预计 15 分钟

Node.js 流式管道实战:用 Transform 流实现 ChatGPT 响应的实时转译与速率控制

一、问题场景

ChatGPT 的 SSE 流式响应是一次 HTTP 请求、多次 data 事件推送。看似简单,但前端到后端的转发链路里会累积三个痛点:

  1. 背压:下游消费慢于上游生产,fetch 的 reader 被迫频繁中断
  2. 粒度过细:每个 token 触发一次 React setState,帧率掉到个位数
  3. 断线丢失:网络抖动导致 SSE 断开,已收到的 token 全废

下面从 Node.js Transform 流切入,把这三个问题逐个拆掉。

二、Transform 流:管道的中间件

Node.js Stream 有 Readable / Writable / Transform / Duplex 四种抽象。Transform 是最适合做"中间加工"的:接收 chunk → 处理 → 吐出结果。

// bad:在可写流里直接 push,语义混乱
const w = new Writable({
  write(chunk, _encoding, cb) {
    this.push(chunk); // ❌ Writable 没有 push 方法
    cb();
  },
});

// good:Transform 天生就是做这个的
const t = new Transform({
  transform(chunk, _encoding, cb) {
    const processed = chunk.toString().toUpperCase();
    this.push(processed);
    cb();
  },
});

每个 chunk 经过 transformthis.push() 放到可读端,cb() 表示处理完成。内置背压:下游读不走时自动暂停上游。

三、实战一:SSE chunk 分割与消息组装

OpenAI 的 SSE 每次推送一条 data: {"choices":[...]}\n\n。但 TCP 流不保证边界,一个 chunk 可能包含半条消息或多条消息。

import { Transform, TransformCallback } from 'node:stream';

class SSEParser extends Transform {
  private buffer = '';

  _transform(chunk: Buffer, _enc: BufferEncoding, cb: TransformCallback) {
    this.buffer += chunk.toString('utf-8');

    const lines = this.buffer.split('\n');
    // 最后一行可能不完整,留到下次
    this.buffer = lines.pop() || '';

    for (const line of lines) {
      if (line.startsWith('data: ')) {
        const data = line.slice(6);
        if (data === '[DONE]') {
          this.push(null); // 结束信号
          return cb();
        }
        this.push(data);
      }
    }
    cb();
  }
}
// bad:直接用 fetch 的 reader,手动处理边界
const reader = response.body!.getReader();
while (true) {
  const { value, done } = await reader.read();
  if (done) break;
  // 自己写 split 和 buffer 逻辑,每次调用都得重新实现
  text += new TextDecoder().decode(value!);
  text.split('\n').filter(l => l.startsWith('data: ')).forEach(/* ... */);
}

// good:管道组合,职责分离
const pipeline = promisify(stream.pipeline);
await pipeline(
  response.body as unknown as Readable,
  new SSEParser(),
  new TokenThrottle({ maxTokensPerSecond: 60 }),
  new BatchMerger({ batchSize: 5, flushInterval: 80 }),
);

四、实战二:速率控制 — 自己写 TokenThrottle

下游太快会丢动画效果,下游太慢用户感到卡顿。好的流式体验要求匀速推进。

class TokenThrottle extends Transform {
  private queue: string[] = [];
  private timer: NodeJS.Timeout | null = null;
  private intervalMs: number;

  constructor(opts: { maxTokensPerSecond: number }) {
    super({ objectMode: true });
    this.intervalMs = 1000 / opts.maxTokensPerSecond;
  }

  _transform(chunk: string, _enc: BufferEncoding, cb: TransformCallback) {
    this.queue.push(chunk);
    this.#drain();
    cb();
  }

  #drain() {
    if (this.timer) return;
    this.timer = setInterval(() => {
      if (this.queue.length === 0) {
        clearInterval(this.timer!);
        this.timer = null;
        return;
      }
      this.push(this.queue.shift());
    }, this.intervalMs);
  }
}

关键点:objectMode: true,每个 push() 输出的不再是 Buffer 而是单个 token 字符串,下游 Transform 直接处理。

五、实战三:批量合并 — 减少 React 渲染次数

每个 token 触发一次 re-render 是性能灾难。这里用一个 BatchMerger 做 5 个 token 或 80ms 的合并窗口。

class BatchMerger extends Transform {
  private batch: string[] = [];
  private timer: ReturnType<typeof setTimeout> | null = null;

  constructor(private opts: { batchSize: number; flushInterval: number }) {
    super({ objectMode: true });
  }

  _transform(chunk: string, _enc: BufferEncoding, cb: TransformCallback) {
    this.batch.push(chunk);
    if (this.batch.length >= this.opts.batchSize) {
      this.#flush();
    } else if (!this.timer) {
      this.timer = setTimeout(() => this.#flush(), this.opts.flushInterval);
    }
    cb();
  }

  _flush(cb: TransformCallback) {
    this.#flush();
    cb();
  }

  #flush() {
    if (this.batch.length === 0) return;
    clearTimeout(this.timer!);
    this.timer = null;
    this.push(this.batch.join(''));
    this.batch = [];
  }
}

_flush 是管道结束时的兜底调用,确保最后不足 5 个的 token 也能输出。

六、React 侧用 useReducer 收敛状态

SSE 的每次管道输出是一个增量字符串。用 useState 叠加会触发依赖旧值的闭包陷阱,useReducer 把状态更新收敛为纯函数:

type State = { content: string; status: 'idle' | 'streaming' | 'done' | 'error' };
type Action =
  | { type: 'APPEND'; payload: string }
  | { type: 'DONE' }
  | { type: 'ERROR'; payload: string };

function reducer(state: State, action: Action): State {
  switch (action.type) {
    case 'APPEND':
      return { ...state, content: state.content + action.payload };
    case 'DONE':
      return { ...state, status: 'done' };
    case 'ERROR':
      return { ...state, status: 'error', content: state.content }; // 保留已收到内容
    default:
      return state;
  }
}
// bad:useState 闭包陷阱
const [content, setContent] = useState('');
source.on('data', (chunk) => {
  setContent(content + chunk); // ❌ content 是闭包捕获的旧值
});

// good:useReducer 每次拿到最新 state
const [state, dispatch] = useReducer(reducer, initialState);
source.on('data', (chunk) => dispatch({ type: 'APPEND', payload: chunk }));

七、断线续传:管道位置快照

Stream 管道本身不记录消费进度。可以加一个 PassThrough 包装器做位点记录:

class CheckpointStream extends Transform {
  public byteCount = 0;

  _transform(chunk: Buffer, enc: BufferEncoding, cb: TransformCallback) {
    this.byteCount += chunk.length;
    this.push(chunk);
    cb();
  }
}

// 断线时 byteCount 就是已接收字节数,下次请求带上:
// fetch('/api/chat', { body: JSON.stringify({ ...prev, offset: checkpoint.byteCount }) })

服务端收到 offset 后跳过已处理的 chunk,利用 TokenThrottle 和 BatchMerger 无缝接续流。


整套方案的依赖关系:TCP chunk → SSEParser → TokenThrottle → BatchMerger → React useReducer。Node.js Stream 的管道模型天然支持背压和组合,比手动 while (reader.read()) 可靠性高一档。下次遇到流式场景,先问自己:这个管道少了一节 Transform,是哪一节?

0 评论

评论区

登录 后参与评论