流与管道
流(Stream)是 Node.js 处理数据最优雅的抽象:把数据切成一块块,边读边处理。文件复制、日志处理、网络传输、压缩解压,都离不开流。理解流,才算真正理解 Node 的 I/O 设计。
一、流的概念与类型
流是"随时间流动的数据集合",数据以 chunk(块) 为单位传递。Node 提供四种流类型:
| 类型 | 方向 | 说明 | 典型示例 |
|---|---|---|---|
Readable | 可读 | 从中读出数据 | fs.createReadStream、http.IncomingMessage |
Writable | 可写 | 向其中写入数据 | fs.createWriteStream、res(响应对象) |
Duplex | 双向 | 既可读又可写,两边独立 | net.Socket、TCP 连接 |
Transform | 转换 | 双向,且读写数据会经过转换 | zlib.createGzip、crypto.createCipher |
javascript
// 判断一个对象属于哪种流
const stream = require("stream");
const fs = require("fs");
const readable = fs.createReadStream("data.txt");
console.log(readable instanceof stream.Readable); // true
console.log(readable instanceof stream.Writable); // false二、读流 Readable
读流有两种消费模式:事件模式(默认)与暂停模式。
javascript
const fs = require("fs");
// 监听 data 事件:进入流动模式,不断吐出数据块
const rs = fs.createReadStream("big.log", { highWaterMark: 16 * 1024 });
rs.on("data", (chunk) => {
console.log("收到一块数据,大小:", chunk.length);
});
rs.on("end", () => {
console.log("读取完毕");
});
// 暂停与恢复
rs.pause(); // 暂停 data 事件
rs.resume(); // 恢复流动(不处理数据时可用它丢弃数据)
// read() 手动拉取(暂停模式下使用)
const rs2 = fs.createReadStream("data.txt");
rs2.on("readable", () => {
let chunk;
while ((chunk = rs2.read()) !== null) {
console.log("手动读取:", chunk.toString());
}
});| 方法/事件 | 说明 |
|---|---|
data 事件 | 每收到一块数据触发,参数是 chunk |
readable 事件 | 有数据可读时触发,配合 read() 使用 |
read(size) | 手动读取指定大小的数据 |
pause() / resume() | 暂停 / 恢复流动模式 |
end 事件 | 数据全部读完时触发 |
三、写流 Writable
javascript
const fs = require("fs");
const ws = fs.createWriteStream("out.txt");
// write 返回 false 表示内部缓冲区已满(触发背压)
const ok = ws.write("第一行\n");
console.log("缓冲区是否未满:", ok); // true 表示可以继续写
ws.write("第二行\n");
ws.end("最后一行\n"); // end 之后不能再 write
// finish 事件:所有数据落盘后触发
ws.on("finish", () => console.log("写入完成"));| 方法/事件 | 说明 |
|---|---|
write(chunk) | 写入一块数据,返回布尔值 |
end([chunk]) | 结束写入,可带最后一块数据 |
finish 事件 | 所有数据已被底层系统接收 |
drain 事件 | 缓冲区清空,可以继续 write |
四、背压机制(Backpressure)
背压:当写端消费速度跟不上读端生产速度时,数据会在缓冲区堆积。Node 用两个机制处理:
- highWaterMark:缓冲区阈值(默认 64KB),超过后
write()返回false - drain 事件:缓冲区低于阈值后触发,此时才能安全继续写入
javascript
const fs = require("fs");
const rs = fs.createReadStream("huge.txt");
const ws = fs.createWriteStream("copy.txt", { highWaterMark: 8 * 1024 });
rs.on("data", (chunk) => {
const canContinue = ws.write(chunk);
// 写不下了,暂停读取,等 drain 再继续
if (!canContinue) {
rs.pause();
ws.once("drain", () => rs.resume());
}
});
rs.on("end", () => ws.end());| 概念 | 含义 |
|---|---|
highWaterMark | 内部缓冲区大小阈值(字节) |
write() 返回 false | 缓冲区已满,触发背压 |
drain 事件 | 缓冲区清空,继续写入的信号 |
| 不加背压的后果 | 内存被缓冲区撑爆、进程崩溃 |
五、pipe 管道
pipe 是背压处理的自动化:readable.pipe(writable) 内部自动完成暂停/恢复,无需手写 drain 逻辑。
javascript
const fs = require("fs");
// 一行代码完成大文件复制,且自动处理背压
fs.createReadStream("big.iso").pipe(fs.createWriteStream("big-copy.iso"));
// pipe 返回目标流,可以继续链式调用
fs.createReadStream("a.txt")
.pipe(fs.createWriteStream("b.txt"))
.on("finish", () => console.log("复制完成"));| 写法 | 优点 | 缺点 |
|---|---|---|
手写 data + write + drain | 可精确控制 | 代码冗长易错 |
pipe | 自动背压、简洁 | 错误不会自动传播(需监听) |
stream/promises 的 pipeline | 自动背压 + 自动错误处理 | 需要 Node 15+ |
javascript
// pipeline 是生产环境推荐写法:错误会自动销毁所有流
const { pipeline } = require("stream/promises");
try {
await pipeline(
fs.createReadStream("big.iso"),
fs.createWriteStream("copy.iso")
);
console.log("复制成功");
} catch (err) {
console.error("复制失败:", err.message);
}六、Transform 流
Transform 流在数据流经时做转换,读到的数据是"转换后"的。
6.1 压缩与解压
javascript
const zlib = require("zlib");
const fs = require("fs");
// 压缩:文件 → gzip 转换流 → 压缩文件
fs.createReadStream("data.txt")
.pipe(zlib.createGzip())
.pipe(fs.createWriteStream("data.txt.gz"));
// 解压
fs.createReadStream("data.txt.gz")
.pipe(zlib.createGunzip())
.pipe(fs.createWriteStream("data-restored.txt"));6.2 加密
javascript
const crypto = require("crypto");
const fs = require("fs");
const key = crypto.randomBytes(32);
const iv = crypto.randomBytes(16);
// 加密写文件
const cipher = crypto.createCipheriv("aes-256-cbc", key, iv);
fs.createReadStream("secret.txt")
.pipe(cipher)
.pipe(fs.createWriteStream("secret.enc"));
// 解密读文件
const decipher = crypto.createDecipheriv("aes-256-cbc", key, iv);
fs.createReadStream("secret.enc")
.pipe(decipher)
.pipe(fs.createWriteStream("secret-out.txt"));6.3 解析 CSV(分块处理)
javascript
const { Transform } = require("stream");
// 自定义转换:把每行 CSV 转成 JSON 行
const csvToJson = new Transform({
transform(chunk, encoding, callback) {
const lines = chunk.toString().trim().split("\n");
const jsonLines = lines.map((line) => {
const [name, age] = line.split(",");
return JSON.stringify({ name, age });
});
callback(null, jsonLines.join("\n") + "\n");
},
});
process.stdin.pipe(csvToJson).pipe(process.stdout);七、链式管道
pipe 可以串联多个 Transform,形成流水线:
javascript
const zlib = require("zlib");
const crypto = require("crypto");
const fs = require("fs");
// 读取 → 压缩 → 加密 → 写入
fs.createReadStream("app.log")
.pipe(zlib.createGzip()) // 压缩
.pipe(crypto.createCipheriv("aes-256-cbc", key, iv)) // 加密
.pipe(fs.createWriteStream("app.log.enc"));
// 数据流向:读流 → Transform → Transform → 写流链式管道的优势是每个环节只关心自己的转换,模块之间完全解耦。
八、流事件总览
| 事件 | 触发时机 | 所属流 |
|---|---|---|
data | 每块数据到达 | Readable |
end | 数据全部读完 | Readable |
error | 出错(读写出错、管道断裂) | 所有流 |
finish | 所有数据写完 | Writable |
close | 流及其底层资源关闭 | 所有流 |
drain | 缓冲区清空可继续写 | Writable |
pipe / unpipe | 被 pipe / 取消 pipe | Writable |
error 事件必须处理——不监听 error,流的错误会直接抛到全局导致进程崩溃:
javascript
const fs = require("fs");
const rs = fs.createReadStream("missing.txt");
rs.on("error", (err) => {
console.error("读流出错:", err.code); // ENOENT
});九、大文件处理:一次性读取 vs 流式
用同样 1GB 的文件对比两种方式:
| 对比项 | 一次性 readFile | 流式处理 |
|---|---|---|
| 内存占用 | 约等于文件大小(1GB+) | 恒定为缓冲区大小(约 64KB) |
| 开始处理的时间 | 等全部读完 | 读到第一块即可处理 |
| 适用文件 | 配置文件、模板 | 日志、视频、压缩包 |
| 进程崩溃风险 | 高(内存溢出) | 低 |
javascript
// 反面示例:一次性读取 1GB 文件
const fs = require("fs");
fs.readFile("huge.iso", (err, data) => {
// 内存瞬间飙升,可能 OOM
});
// 正面示例:流式,每块 64KB 依次处理
fs.createReadStream("huge.iso", { highWaterMark: 64 * 1024 })
.on("data", (chunk) => {
// 逐块处理,内存始终平稳
});十、自定义流
通过继承 stream.Readable 实现一个数字生成器:
javascript
const { Readable } = require("stream");
class NumberStream extends Readable {
constructor(max = 10) {
super();
this.max = max;
this.current = 0;
}
// _read 必须实现:调用 push 把数据放进流里
_read() {
if (this.current >= this.max) {
this.push(null); // null 表示结束
return;
}
this.push(`${this.current}\n`);
this.current += 1;
}
}
const nums = new NumberStream(5);
nums.on("data", (chunk) => process.stdout.write(chunk));
// 输出 0 1 2 3 4自定义 Transform 流实现大小写转换:
javascript
const { Transform } = require("stream");
class UpperCaseStream extends Transform {
_transform(chunk, encoding, callback) {
this.push(chunk.toString().toUpperCase());
callback();
}
}
process.stdin.pipe(new UpperCaseStream()).pipe(process.stdout);| 要实现的方法 | 所在流 | 作用 |
|---|---|---|
_read() | Readable | 通过 push(data) 产出数据 |
_write(chunk, enc, cb) | Writable | 处理写入的数据,调用 cb() 通知完成 |
_transform(chunk, enc, cb) | Transform | 转换数据,this.push() 输出 |
流让数据像水管一样流动:源头是读流,出口是写流,中间可以串接任意多个转换环节。掌握背压与管道,就能稳健地处理任意大小的数据。