流
Node.js 流(Stream)详解
流是 Node.js 中处理有序数据的抽象接口,用于以连续方式异步读取或写入数据。它们特别适合处理大文件、网络通信或任何可能无限产生的数据。Node.js 内置的 stream 模块提供了四种基本流类型,以及用于组合和错误处理的工具函数。
一、流的四种基本类型
| 类型 | 描述 | 方向 | 典型应用 |
|---|---|---|---|
| 可读流 (Readable) | 数据可以从中读取 | 单向:流出 | fs.createReadStream、http.IncomingMessage |
| 可写流 (Writable) | 数据可以写入其中 | 单向:流入 | fs.createWriteStream、http.ServerResponse |
| 双工流 (Duplex) | 既可读又可写 | 双向独立 | net.Socket、crypto.createCipheriv |
| 转换流 (Transform) | 双工流的一种,会对数据进行转换 | 双向且数据经过处理 | zlib.createGzip、crypto.createHash |
示例代码
const { Readable, Writable, Duplex, Transform } = require("stream");
// 可读流
const readable = Readable.from(["Hello", "World"]);
// 可写流
const writable = fs.createWriteStream("./output.txt");
// 双工流
const duplex = new Duplex({
read(size) {
/* 实现读 */
},
write(chunk, encoding, callback) {
/* 实现写 */
},
});
// 转换流(大小写转换)
const transform = new Transform({
transform(chunk, encoding, callback) {
const output = chunk.toString().toUpperCase();
callback(null, output);
},
});
二、背压机制与 pipe / pipeline
1. 背压机制(Backpressure)
当可读流的数据产生速度超过可写流的消费速度时,数据会在内部缓冲区累积。背压机制通过控制数据流动来防止内存溢出:
- 可读流有一个内部
highWaterMark(默认 16KB)。当缓冲区大小超过该值,readable.push()返回false,表示应暂停读取。 - 可写流也有
highWaterMark,当写入速度慢时,writable.write()返回false,并触发drain事件。 - 在
pipe或pipeline中,背压自动处理;手动实现时需要检查返回值并监听drain。
2. readable.pipe(writable)
最简单的管道方式,自动处理背压和结束事件。缺点:错误处理不完善(需要单独监听 error 事件)。
const readStream = fs.createReadStream("bigfile.txt");
const writeStream = fs.createWriteStream("copy.txt");
readStream.pipe(writeStream);
writeStream.on("finish", () => console.log("复制完成"));
3. stream.pipeline(推荐)
从 Node.js 10 开始引入,支持 Promise 并自动清理流、转发错误,避免内存泄漏。
const { pipeline } = require("stream");
const { promisify } = require("util");
const pipelineAsync = promisify(pipeline);
// 或者使用 stream/promises 直接导入 Promise 版本
const { pipeline: pipelinePromise } = require("stream/promises");
async function run() {
await pipelinePromise(
fs.createReadStream("input.txt"),
zlib.createGzip(),
fs.createWriteStream("input.gz")
);
console.log("压缩完成");
}
三、自定义流实现
通过继承核心流类并实现必要的方法,可以创建定制流。
1. 自定义可读流
需要实现 _read(size) 方法,内部调用 push(chunk) 或 push(null) 结束。
const { Readable } = require("stream");
class Counter extends Readable {
constructor(options) {
super(options);
this.max = 10;
this.index = 0;
}
_read(size) {
if (this.index < this.max) {
this.push(String(this.index++));
} else {
this.push(null); // 结束
}
}
}
const counter = new Counter();
counter.pipe(process.stdout); // 输出 0 1 2 ... 9
2. 自定义可写流
需要实现 _write(chunk, encoding, callback),处理数据后调用 callback()。
const { Writable } = require("stream");
class MyWritable extends Writable {
_write(chunk, encoding, callback) {
console.log(`写入: ${chunk.toString()}`);
callback(); // 继续接收下一个块
}
}
const writable = new MyWritable();
writable.write("Hello\n");
writable.write("World");
writable.end();
3. 自定义转换流
实现 _transform(chunk, encoding, callback),处理后调用 callback(null, transformedChunk)。
const { Transform } = require("stream");
class UpperCaseTransform extends Transform {
_transform(chunk, encoding, callback) {
const upper = chunk.toString().toUpperCase();
callback(null, upper);
}
}
const transform = new UpperCaseTransform();
process.stdin.pipe(transform).pipe(process.stdout);
4. 自定义双工流
需同时实现 _read 和 _write,且两者独立。如果不涉及数据转换,直接使用 Duplex。
四、实用工具:stream/promises 中的 pipeline 和 finished
1. pipeline(Promise 版本)
用于串联多个流,返回 Promise,并在任何流发生错误时中止并清理。
const { pipeline } = require("stream/promises");
const fs = require("fs");
const zlib = require("zlib");
async function compress() {
await pipeline(
fs.createReadStream("file.txt"),
zlib.createGzip(),
fs.createWriteStream("file.gz")
);
console.log("压缩成功");
}
2. finished
检测流是否完成(或出错),返回 Promise。常用于手动管理流结束。
const { finished } = require("stream/promises");
async function monitor(stream) {
try {
await finished(stream);
console.log("流已完成");
} catch (err) {
console.error("流错误:", err);
}
}
const rs = fs.createReadStream("somefile.txt");
monitor(rs);
finished 也可用于每个独立流,而 pipeline 更适合多流串联。
注意:
stream/promises从 Node.js 15 开始稳定可用,对于较低版本可使用util.promisify(require('stream').finished)。
五、背压手动处理示例
当不使用 pipe 或 pipeline 时,需要手动处理背压:
const rs = fs.createReadStream("large.txt");
const ws = fs.createWriteStream("copy.txt");
rs.on("data", chunk => {
const canWrite = ws.write(chunk);
if (!canWrite) {
rs.pause(); // 暂停读取
ws.once("drain", () => rs.resume()); // 等待可写流排空
}
});
rs.on("end", () => ws.end());
ws.on("finish", () => console.log("复制结束"));
pipe 内部就是实现这样的逻辑。
六、对象模式(Object Mode)
默认流传输的是 Buffer 或字符串,通过设置 objectMode: true 可以传输任意 JavaScript 对象(如 JSON 数据),但会禁用背压控制,需谨慎使用。
const { Readable } = require("stream");
const objStream = new Readable({
objectMode: true,
read() {
this.push({ id: 1 });
this.push({ id: 2 });
this.push(null);
},
});
objStream.on("data", obj => console.log(obj));
总结
| 流类型 | 核心方法 | 用途 |
|---|---|---|
Readable | _read | 产生数据 |
Writable | _write | 消费数据 |
Duplex | _read + _write | 双向通信 |
Transform | _transform | 数据转换(如压缩、加密) |
关键最佳实践:
- 使用
pipeline组合流,自动处理背压和错误。 - 避免直接使用
pipe而不监听错误。 - 自定义流时继承正确的基类,并实现必需方法。
- 使用
finished检测流结束状态。
推荐参考资料
- Node.js 官方文档:Stream API
- Backpressuring in Streams(背压指南)
- 《Node.js 设计模式》第 4 章:流式编程
相关文章
文件系统
Node.js 的 fs 模块提供了与文件系统交互的 API,几乎涵盖了所有标准文件操作。它支持三种风格的 API:同步、回调式异步 和 Promise 式异步。下面详细介绍这些 API 的选择策略、流式读写、文件监视以及常用操作。
缓冲区
Buffer 是 Node.js 全局对象,用于处理二进制数据流(如文件、网络数据)。在 ES6 引入 TypedArray 之前,Buffer 是 Node.js 处理二进制的主要方式;现在 Buffer 实现了 Uint8Arra…
路径处理
path 模块提供了用于处理和转换文件路径的实用工具。它是 Node.js 核心模块,无需安装即可直接使用。由于不同操作系统(Windows、Linux、macOS)的路径分隔符不同(Windows 使用反斜杠 \,POSIX 使用正…
进程与子进程
Node.js 在单个进程中运行,但通过 process 对象、child_process 模块及 worker_threads 模块提供了对进程的精细控制以及多进程/多线程能力。下面从三个方面展开。
网络编程
Node.js 提供了丰富的内置模块,用于构建网络应用。下面从 HTTP/HTTPS、WebSocket、TCP/UDP 以及 DNS 解析四个方面进行深入介绍。
其他常用模块
以下是 Node.js 中几个重要内置模块的详细介绍,涵盖事件触发器、定时器、加密、压缩、操作系统信息和实用工具。
Series
core
7 / 7