姆姆极客分享MU·GEEK·SHARE
// Node.js Stream Pipeline Readable fs / HTTP Transform async fn* Writable res / file highWaterMark backpressure drain pipeline() · auto backpressure · auto cleanup
后端架构·

Node.js 流式处理与背压机制深度实践:从 fs 到 web stream

Stream 是 Node.js 最被低估的能力之一。本文从 backpressure 原理出发,对比 classic stream 与 web stream,演示如何写出"既快又不爆内存"的流式服务。


Node.js 的 stream API 是历史最久、也被吐槽最多的 API 之一。但吐槽归吐槽,stream 在文件处理、HTTP 代理、日志聚合、ETL 这些场景里依然是不可替代的工具。本文从 backpressure 这个核心概念讲起,把 stream 写“对”的几个关键点讲清楚。

一、为什么需要 backpressure

想象一个场景:你要把一个 10GB 的文件通过 HTTP 传给客户端。如果直接 fs.readFile 然后返回,内存会被吃满;如果用 stream 但不处理背压,磁盘读得快、网络发得慢,buffer 会在内存里堆到 OOM。

// 反例:忽略了背压
const readable = fs.createReadStream('huge.txt');
const writable = res;  // HTTP 响应

readable.on('data', (chunk) => {
  writable.write(chunk);  // ← 网络慢时这里会无限堆积
});

writable.write() 返回 false 表示“内部 buffer 已满,请等我 drain”。但上面这段代码完全忽略了返回值,所有 chunk 都被强塞进去——内部 buffer 满了就转到底层 nodejs buffer,直到 OOM。

二、正确姿势:用 pipeline

import { pipeline } from 'node:stream/promises';
import { createReadStream, createWriteStream } from 'node:fs';

await pipeline(
  createReadStream('input.txt'),
  // 中间可以插入任意 Transform
  async function* (source) {
    for await (const chunk of source) {
      yield chunk.toString().toUpperCase();
    }
  },
  createWriteStream('output.txt'),
);

pipeline 自动处理背压、错误传递、资源清理。它是 stream 处理的官方推荐方式,应该完全替代手写的 pipe()data 事件监听。

三、stream 的三种类型与生成器

Node.js stream 有四种基础类型(Readable / Writable / Duplex / Transform),但实际开发中你更常用的是生成器函数

import { Readable } from 'node:stream';

// 从任意异步源创建 Readable
const stream = Readable.from((async function* () {
  for (let i = 0; i < 1000; i++) {
    await new Promise(r => setTimeout(r, 10));
    yield { id: i, data: `item-${i}` };
  }
})());

for await (const item of stream) {
  console.log(item);
}

生成器天然支持 for await...of,且自动集成背压——下游消费慢时,生成器会自动暂停在 yield

四、Web Streams:未来已来

Node.js 从 16.5 开始支持 Web Streams API(ReadableStream / WritableStream / TransformStream),这是 WHATWG 标准化的 stream,与浏览器、Deno、Bun 通用。

// Node.js + fetch(返回的是 web ReadableStream)
const res = await fetch('https://api.example.com/large');
const reader = res.body.getReader();

while (true) {
  const { done, value } = await reader.read();
  if (done) break;
  // value 是 Uint8Array
  await process(value);
}

Web Streams 的核心区别:

  • 没有 data 事件,统一用 getReader()pipeTo()
  • 默认 pull-based,背压机制更自然
  • 跨运行时通用,更适合库代码

互转:fromWeb / toWeb

import { Readable } from 'node:stream';

// 把 Node Readable 转成 Web ReadableStream
const webStream = Readable.toWeb(nodeReadable);

// 反过来
const nodeReadable = Readable.fromWeb(webStream);

写库代码时,优先暴露 Web Streams API,Node.js 用户用 fromWeb 转一下,浏览器用户直接用。这是 2026 年的正确姿势。

五、实战:流式 JSON 解析

传统 JSON.parse 要求完整字符串,对大文件不友好。配合 stream + 增量解析器可以做到边读边解析:

import { pipeline } from 'node:stream/promises';
import { createReadStream } from 'node:fs';
import { parse } from 'stream-json/jsonl/Parser.js';  // stream-json 库

await pipeline(
  createReadStream('large.jsonl'),
  parse(),
  async function* (source) {
    for await (const { value } of source) {
      yield transform(value);  // 增量处理每一行
    }
  },
  // 写入数据库或下游
);

这种“流式 ETL”的内存占用恒定(只与单个 chunk 大小相关),与文件大小无关。

六、HTTP 服务里的流式响应

import { createReadStream } from 'node:fs';
import { pipeline } from 'node:stream/promises';

http.createServer(async (req, res) => {
  res.setHeader('Content-Type', 'video/mp4');

  // 不要 await!流式响应应该让 pipeline 后台跑
  pipeline(
    createReadStream('movie.mp4'),
    res,
  ).catch(err => {
    console.error('stream failed', err);
    if (!res.headersSent) res.statusCode = 500;
    res.end();
  });
});

关键点:

  • 不要 await pipeline,否则请求会等到文件全部传完才返回
  • 监听客户端断开(req.on('close')),pipeline 会自动取消
  • HTTP/2 + Server Push 时这种模式仍然有效

七、Backpressure 与 highWaterMark

每个 stream 都有 highWaterMark 参数,表示“内部 buffer 满到什么程度就触发背压”:

const stream = createReadStream('file.txt', {
  highWaterMark: 64 * 1024,  // 64KB,默认 16KB
});

调大的 trade-off:

  • ✅ 系统调用次数减少(每次 read 更多)
  • ✅ 突发写入更平滑
  • ❌ 内存占用上升
  • ❌ 首字节延迟略增

经验值:

  • 文件 IO:64KB - 1MB
  • 网络 IO:保持默认 16KB
  • 对象模式(objectMode: true):16 - 100 个对象

八、错误处理:pipeline vs pipe

// 反例:pipe 不处理错误
readable.pipe(transform).pipe(writable);
// transform 出错时,readable 不会自动销毁,造成资源泄漏

// 正例:pipeline 自动传播错误 + 清理资源
await pipeline(readable, transform, writable);

pipeline 还有一个杀手锏:任一 stream 出错,它会自动销毁整条链上的所有 stream,避免 fd 泄漏。这是 pipe 永远做不到的。

九、诊断:stream 为什么慢

import { performance, createHistogram } from 'node:perf_hooks';

// 观察 stream 的内部状态
readable.on('data', (chunk) => {
  console.log({
    length: readable.readableLength,  // 当前 buffer 字节数
    highWaterMark: readable.readableHighWaterMark,
    flowing: readable.readableFlowing,
  });
});

如果 readableLength 一直接近 highWaterMark,说明下游消费慢;如果一直接近 0,说明上游生产慢。这是定位 stream 瓶颈的最快方法。

十、性能对比:内存与吞吐

我们做的一个日志聚合服务,从“全量加载”切到“流式处理”的对比:

指标 全量加载 流式处理
处理 1GB 日志的内存峰值 2.1GB 38MB
处理耗时 12s 9s
并发上限(同机) 3 25

结语

Node.js 的 stream API 因为历史包袱被严重低估,但它依然是后端 IO 密集场景最锋利的工具。掌握三个要点就能避免大部分坑:

  1. 永远用 pipeline,不用 pipe
  2. 用生成器函数写自定义流,不要继承 Readable 类
  3. 库代码暴露 Web Streams,应用代码用 Node Streams

当你把这三条变成肌肉记忆,就能写出“快、稳、内存可控”的流式服务——这是后端工程师的基本功,也是面试时的常见区分点。