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
使用 readable 和 read() 主动拉取:
input.on("readable", () => {
let chunk: Buffer | null;
while ((chunk = input.read()) !== null) {
console.log(chunk.length);
}
});
不要同时混用 data、readable、pipe() 和异步迭代,否则数据流动模式和消费归属会变得难以推断。
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() 在此用于丢弃并消费数据。
练习题
- 为什么
close不能单独证明读取成功? - 文件不存在时预计出现哪些事件,哪些不会出现?
- 比较
data与for await...of的代码和取消方式。 - 修改
highWaterMark,观察 Chunk 数量而非假设 Chunk 恒定大小。