Node.js流式数据处理是通过Stream模块以连续分块的形式处理输入与输出,适合大文件读写、网络请求转发和实时日志处理等场景。相比一次性读取全部内容,流能降低内存压力并提高吞吐效率。

Node.js中的流类型
Node.js内置了四种基础流,分别承担不同的数据处理角色:
- Readable:可读流,数据源端,如文件读取流
- Writable:可写流,数据终点,如文件写入流
- Duplex:双工流,可读也可写,如Socket
- Transform:转换流,处理中间数据,如压缩流
使用管道连接流
最常用的方法是调用pipe方法,将可读流的数据自动推送到可写流,不需要手动监听data事件。
const fs = require('fs');
// 创建可读流与可写流
const readable = fs.createReadStream('input.txt');
const writable = fs.createWriteStream('output.txt');
// 使用pipe串联,自动处理背压
readable.pipe(writable);
writable.on('finish', () => {
console.log('文件复制完成');
});
自定义可读流示例
当内置流不能满足需求时,可以继承Readable来推送自定义数据块。
const { Readable } = require('stream');
class MyReadable extends Readable {
constructor(data) {
super();
this.data = data;
this.index = 0;
}
_read() {
if (this.index < this.data.length) {
// 推送一个数据块
this.push(this.data[this.index]);
this.index++;
} else {
// 结束流
this.push(null);
}
}
}
const stream = new MyReadable(['块1', '块2', '块3']);
stream.on('data', (chunk) => {
console.log('收到:', chunk);
});
处理背压问题
当可写流消费速度慢于可读流生产速度时,会产生背压。pipe方法内部已处理该情况,若手动监听data事件则需调用pause与resume控制节奏。
| 方式 | 优点 | 注意点 |
|---|---|---|
| pipe | 自动背压,代码简洁 | 不易插入自定义转换 |
| 手动data | 控制精细 | 需自行处理pause和resume |
转换流的实际运用
Transform流适合在传输过程中修改数据,例如将输入文本转为大写。
const { Transform } = require('stream');
const upper = new Transform({
transform(chunk, encoding, callback) {
// 将缓冲区内容转为大写后推送
this.push(chunk.toString().toUpperCase());
callback();
}
});
process.stdin.pipe(upper).pipe(process.stdout);
通过上述方式,Node.js流式数据处理能够在多种业务里平稳运行,不必担心大体积内容导致进程崩溃。