11 Backpressure:数据为什么不能无限向前推

1. 生产者和消费者速度不一致

假设磁盘读取 500 MiB/s,网络只能发送 10 MiB/s。如果读取端持续把所有数据推入内存,内存会不断增长。背压是消费者告诉生产者:

我暂时处理不过来,请放慢。

2. 手写复制的错误示例

input.on("data", (chunk) => {
  output.write(chunk);
});

它忽略 write() 返回值。正确的手工方向是:

input.on("data", (chunk) => {
  if (!output.write(chunk)) {
    input.pause();
  }
});

output.on("drain", () => {
  input.resume();
});

但还没有处理双方 errorclose、结束和销毁,所以生产代码仍应使用 pipeline()

3. 缓冲在哪里

多级链路中每一段都可能有缓冲:

File ReadStream buffer
  → Markdown/Transform buffer
  → gzip buffer
  → HTTP Response/socket buffer
  → 操作系统内核 buffer

highWaterMark: 64 KiB 只描述某个 Stream 的阈值,不代表整个请求最多使用 64 KiB。

4. pipe() 如何协助

readable.pipe(writable) 会根据目标 Writable 的 write() 返回值自动暂停/恢复 Readable,解决基本背压。但它对多段错误和资源清理的保证不如 pipeline() 完整。

5. Object Mode

对象模式的 highWaterMark 计量对象数量,不是字节数:

new Transform({
  objectMode: true,
  highWaterMark: 16,
});

16 个巨大对象仍然可能占用大量内存。把整份表达矩阵装进一个对象后传递,不能因为使用 Stream 就认为安全。

6. ZIP 转换中的最慢阶段

ZIP Entry 读取很快
Markdown 渲染较慢
磁盘写入一般
Archiver 压缩 CPU 较慢
客户端下载更慢

只要某一步没有正确传递背压,上游就可能积压。尤其不要在 data 事件中启动大量未等待的异步任务:

input.on("data", async (chunk) => {
  await slowProcess(chunk);
});

EventEmitter 不会等待异步监听器,后续 Chunk 仍会继续到来。

7. 观察背压

创建一个延迟写入的 Writable,记录:

write() 何时返回 false
drain 何时发生
RSS 如何变化
输入是否暂停

不要用固定 sleep() 代替背压协议;等待 drain 才表示目标重新有缓冲能力。

练习题

  1. 为什么 data 的 async 监听器不会自然形成背压?
  2. 多级 Pipeline 中有多少处可能缓冲?
  3. Object Mode 的 highWaterMark: 16 为什么不等于 16 字节?
  4. 分别实现错误版和正确版复制,对比 RSS。