09 Node.js Readable Stream

1. Readable 是“可被消费的数据源”

文件、HTTP 请求、ZIP Entry 都可能是 Node.js Readable

import { createReadStream } from "node:fs";

const input = createReadStream(filePath, {
  highWaterMark: 64 * 1024,
});

它不会把整个文件一次性加载到内存,而是在底层读取与消费者速度之间维护内部缓冲区。

2. 两种主要模式

Flowing mode

注册 data 监听器或调用 pipe() 后,数据主动流出:

input.on("data", (chunk: Buffer) => {
  console.log(chunk.length);
});

Paused/readable mode

使用 readableread() 主动拉取:

input.on("readable", () => {
  let chunk: Buffer | null;

  while ((chunk = input.read()) !== null) {
    console.log(chunk.length);
  }
});

不要同时混用 datareadablepipe() 和异步迭代,否则数据流动模式和消费归属会变得难以推断。

3. 推荐消费方式

需要自己处理 Chunk 时,可用异步迭代:

let totalBytes = 0;

for await (const chunk of input) {
  totalBytes += (chunk as Buffer).length;
}

需要连接到目标 Stream 时,优先使用 pipeline(),不要手工收集所有 Chunk。

4. 关键事件

事件 含义 是否等于成功完成
open 文件描述符已打开,仅文件 Stream
ready 文件 Stream 可以使用
data 一个 Chunk 被消费
readable 缓冲中可能有数据可读
end 所有数据已经被消费者读完 是“读取结束”,不是资源关闭
error 读取或底层资源失败
close Stream/底层资源关闭 不说明成功还是失败

aborted 主要是 http.IncomingMessage 的语义,不是所有通用 Readable 都会发出。

5. 正常时间线

文件 Stream 常见顺序:

open
  → ready
  → data ...
  → end
  → close

不要把它当成所有 Readable 的绝对顺序:内存 Readable 没有文件句柄,也可能没有 open

6. 错误时间线

open
  → data ...
  → error
  → close(默认 emitClose 时通常出现)

一旦读取失败,end 通常不会出现。因此下面的 Promise 可能永久等待:

await new Promise<void>((resolve) => {
  input.once("end", resolve);
});

它没有处理 error。更好的方式是 pipeline(),或在只观察单个 Stream 时使用 finished()

7. pause() 不等于取消

input.pause();
input.resume();

暂停只是暂时停止 data 流出,文件句柄和缓冲数据仍然存在。取消应使用:

input.destroy(new Error("任务取消"));

destroy() 会让 Stream 进入销毁流程,之后不能当成新的数据源复用。

8. highWaterMark

highWaterMark 是开始施加背压的阈值提示,不是严格的总内存上限,也不是每次 Chunk 必然大小。多级 Stream、并发任务和库内部缓冲都会增加总内存。

9. 状态属性

排查时可以观察:

console.log({
  destroyed: input.destroyed,
  readableEnded: input.readableEnded,
  readableFlowing: input.readableFlowing,
});

状态用于调试,不应靠轮询状态替代事件、pipeline()finished()

实操:记录事件顺序

for (const eventName of [
  "open",
  "ready",
  "end",
  "close",
  "error",
] as const) {
  input.on(eventName, (...args: unknown[]) => {
    console.log(Date.now(), eventName, args[0]);
  });
}

input.resume();

没有消费者时,end 不一定立即出现;resume() 在此用于丢弃并消费数据。

练习题

  1. 为什么 close 不能单独证明读取成功?
  2. 文件不存在时预计出现哪些事件,哪些不会出现?
  3. 比较 datafor await...of 的代码和取消方式。
  4. 修改 highWaterMark,观察 Chunk 数量而非假设 Chunk 恒定大小。