上一节我们强调了流的核心价值——分块处理大文件、降低内存占用、通过背压自动调节数据生产与消费速度。Node.js 在实现这一机制时,将流划分为四种基本类型,每种类型都有清晰的职责和典型的使用场景。理解这四种流,是灵活运用 pipe、pipeline 以及自定义流的基础。
7.2.1 可读流(Readable Stream)
可读流是数据的生产者,它负责从某个数据源(文件、网络、内存等)中读取数据,并以流的形式输出给消费者。
基本工作原理
可读流内部维护了一个缓冲区,数据从底层的源被推入缓冲区,消费者通过监听 data 事件或者调用 read() 方法来获取数据。流有两种工作模式:
- 流动模式(flowing mode):数据自动从底层源读出,并通过
data事件尽可能快地推送给消费者。只需添加data事件监听器,流就会切换到该模式。 - 暂停模式(paused mode):必须显式调用
stream.read()方法来拉取数据。这是可读流的初始模式,添加readable事件监听器会保持此模式。
对于大多数场景,直接用 data 事件或管道 pipe 会更方便,它们会自动处理模式切换。
常见可读流实例
// 1. 文件可读流:分块读取大文件
const fs = require('fs');
const readStream = fs.createReadStream('./large-file.txt', { highWaterMark: 64 * 1024 }); // 每次读 64KB
readStream.on('data', (chunk) => {
console.log(`收到 ${chunk.length} 字节数据`);
});
readStream.on('end', () => {
console.log('文件读取完毕');
});
// 2. HTTP 请求的 req 对象也是一个可读流
const http = require('http');
http.createServer((req, res) => {
// req 是一个可读流,包含了客户端发送的请求体数据
req.on('data', (chunk) => {
// 处理请求体块
});
req.on('end', () => {
res.end('接收完毕');
});
}).listen(3000);
// 3. process.stdin 标准输入也是可读流
process.stdin.on('data', (chunk) => {
console.log(`输入: ${chunk.toString()}`);
});
关键事件和方法
data事件:当流将数据块的所有权传递给消费者时触发。end事件:当流中没有更多数据可供消费时触发。error事件:读取过程中出现错误时触发(务必监听,否则错误可能使进程崩溃)。pipe()方法:将可读流连接到可写流,自动处理背压和结束事件。pause()/resume():手动切换流动/暂停模式。read(size):主动从内部缓冲区读取指定大小的数据。
背压处理
当可读流的读取速度远快于可写流的写入速度时,pipe 或 pipeline 会自动暂停可读流的数据输出,直到缓冲区被消耗到一定程度,这就是背压机制的体现。开发者通常不需要手动干预,除非自己实现自定义的流处理逻辑。
7.2.2 可写流(Writable Stream)
可写流是数据的消费者,它负责将收到的数据写入目标(文件、网络套接字、HTTP 响应等)。
基本工作原理
可写流内部也有一个缓冲区,write() 方法会将数据放在这个缓冲区中排队,然后异步地写入底层目标。当缓冲区已满(达到 highWaterMark 阈值),write() 会返回 false,通知生产者暂停写入,直到 drain 事件触发,表示缓冲区已排空,可以继续写入。
常见可写流实例
// 1. 文件可写流:逐块写入大文件
const fs = require('fs');
const writeStream = fs.createWriteStream('./output.txt');
writeStream.write('第一行数据\n');
writeStream.write('第二行数据\n');
writeStream.end('最后一行'); // 写入最后数据并关闭流
writeStream.on('finish', () => {
console.log('文件写入完成');
});
// 2. HTTP 响应的 res 对象是一个可写流
http.createServer((req, res) => {
res.writeHead(200, { 'Content-Type': 'text/plain' });
res.write('Hello ');
res.write('World');
res.end();
});
// 3. process.stdout 标准输出是可写流
process.stdout.write('这是一条输出\n');
关键事件和方法
write(chunk, [encoding], [callback]):写入数据块,返回布尔值指示是否需要等待drain事件。end([chunk], [encoding], [callback]):结束流,可选地写入最后一块数据。调用后不能再写入。finish事件:所有数据已被写入到底层系统,并且end()被调用后触发。error事件:写入或管道中出现错误时触发。drain事件:当内部缓冲区排空,可以继续安全写入时触发,是实现背压控制的关键。
背压控制示例
在手动调用 write() 时,必须处理返回值:
function writeData(writeStream, data, callback) {
let i = 0;
function write() {
let ok = true;
while (i < data.length && ok) {
const chunk = data[i];
i++;
if (i === data.length) {
// 最后一块,使用 end 结束
writeStream.end(chunk, callback);
} else {
ok = writeStream.write(chunk);
}
}
if (i < data.length) {
// 缓冲区已满,等待 drain 后继续
writeStream.once('drain', write);
}
}
write();
}
7.2.3 双工流(Duplex Stream)
双工流同时实现了可读流和可写流接口,它既是一个数据生产者,也是一个数据消费者。它的内部通常维护着独立的输入缓冲区和输出缓冲区,读写操作互不干扰。
典型场景与实例
双工流最常见的例子是网络套接字,如 TCP 连接:
const net = require('net');
const server = net.createServer((socket) => {
// socket 就是一个双工流
socket.on('data', (chunk) => {
console.log(`收到客户端数据: ${chunk.toString()}`);
// 可以向同一个 socket 写回数据
socket.write(`服务器应答: ${chunk.toString()}`);
});
socket.on('end', () => {
console.log('客户端断开连接');
});
});
server.listen(8080);
在这个例子中,socket 既可以通过 write 方法发送数据给客户端(可写流),又可以通过监听 data 事件接收客户端发来的数据(可读流)。两个方向的通信完全独立。
其他双工流实例包括:
zlib.createDeflate等压缩流实际上就是双工流(同时也是转换流)。- 加密流
crypto.createCipheriv:接收明文、输出密文,但同时需要传入密钥等初始化参数。 - 自定义双工流:通过继承
stream.Duplex并实现_read和_write方法。
双工流与转换流的区别
很多开发者容易混淆双工流和转换流。它们的核心区别在于:
- 在双工流中,读和写是完全解耦的,写入的数据和读出的数据之间不一定存在变换关系。例如一个 TCP 套接字,你写入什么内容,读取时可能收到完全不同的内容。
- 转换流(Transform)对数据进行了某种变换,写入的内容经过处理后成为读取的内容,输入和输出存在因果关系。
可以理解为:所有转换流都是双工流,但双工流不一定是转换流。
7.2.4 转换流(Transform Stream)
转换流是一种特殊的双工流,它会将写入端的数据经过某种处理后,从读取端输出。 换句话说,转换流在内部完成了“读取-变换-输出”的串联,是数据管道中的中间处理节点,非常适合做数据格式转换、压缩、加密等工作。
基本原理
转换流内部自动连接了 _transform 方法。当数据块被写入时,会调用 _transform(chunk, encoding, callback),在这个方法里可以对数据块进行修改、缓存或分块,然后通过 callback 将变换后的块推入可读端。此外,还可以实现 _flush(callback) 方法,在流结束前处理剩余的数据。
典型转换流实例
const { Transform } = require('stream');
// 创建一个将输入转换为大写并加上行号的转换流
const lineTransform = new Transform({
transform(chunk, encoding, callback) {
// chunk 是 Buffer,先转为字符串处理
const lines = chunk.toString().split('\n');
const transformed = lines
.filter(line => line) // 过滤空行
.map((line, index) => `${index + 1}: ${line.toUpperCase()}`)
.join('\n') + '\n';
callback(null, transformed);
}
});
// 使用管道:标准输入 -> 转换流 -> 标准输出
process.stdin.pipe(lineTransform).pipe(process.stdout);
在这个例子中,终端输入的内容经过 lineTransform 处理后变成大写并编号,然后显示在标准输出。对开发者来说,只是连接了几个流,数据就在流通中完成了转换。
Node.js 内置的转换流应用非常广泛:
zlib.createGzip()/zlib.createGunzip():压缩和解压缩流,用于 HTTP 压缩响应、文件压缩等。crypto.createCipheriv()/crypto.createDecipheriv():加密和解密流,用于数据保密传输。stream.PassThrough:一个不做任何变换的转换流,通常用于测试或作为占位符。readline.createInterface({ input, output }):虽然不是直接的流,但基于转换原理。
自定义转换流与 backpressure 自动处理
当我们通过 pipe 连接可读流、转换流和可写流时,背压自动贯穿整条管道。如果可写流返回 false,转换流的 _transform 方法会被暂停调用,直到 drain 事件触发。这意味着你的自定义转换流无需额外处理即可获得完整的背压保护。
一个更实际的例子,实现一个简单的 JSON 行解析器转换流:
const { Transform } = require('stream');
class JsonToCsvTransform extends Transform {
constructor() {
super({ readableObjectMode: true, writableObjectMode: false });
this.buffer = '';
}
_transform(chunk, encoding, callback) {
this.buffer += chunk.toString();
const lines = this.buffer.split('\n');
// 最后一个可能是不完整的行,保留到下次处理
this.buffer = lines.pop();
for (const line of lines) {
if (line.trim()) {
try {
const obj = JSON.parse(line);
// 假设 obj 是 { name, age },转为 CSv 行
this.push(`${obj.name},${obj.age}\n`);
} catch (err) {
callback(err);
return;
}
}
}
callback();
}
_flush(callback) {
// 处理最后剩余的数据
if (this.buffer.trim()) {
try {
const obj = JSON.parse(this.buffer);
this.push(`${obj.name},${obj.age}\n`);
} catch (err) {
callback(err);
return;
}
}
callback();
}
}
// 使用:读取 JSON 行文件,转换为 CSV 写入
fs.createReadStream('input.jsonl')
.pipe(new JsonToCsvTransform())
.pipe(fs.createWriteStream('output.csv'));
7.2.5 四种流类型的关系总结
用一张简单的关系图可以概括:
可读流(Readable) --------> 转换流(Transform) --------> 可写流(Writable)
| ↑ ↑ ↑
| | | |
+-------- 双工流(Duplex) --+ +-------------------------+
- 可读流:仅产生数据,不可写入。
- 可写流:仅消费数据,不可读取。
- 双工流:既可读又可写,且读写无直接变换关系。
- 转换流:是双工流的子集,写入的内容经过变换后成为可读取的内容,读写之间存在逻辑关联。
掌握这四种流类型,就掌握了 Node.js 流编程的核心词汇。在后续小节中,我们将深入背压的工作原理和 pipeline 的健壮用法,进一步展示如何将这些流类型组合成高效、健壮的数据处理管道。