Skip to content

Node.js 流 (Streams)

流是 Node.js 中的基本概念之一,提供了一种高效的方式来按顺序处理读写数据。将它们想象成数据的传送带:流允许你分块处理数据,而不是像使用 fs.readFile 那样等待整个数据集加载到内存中。当数据到达或准备好发送时,就可以处理这些小块数据。

这使得流对于处理大量数据(如大文件或网络响应)特别有用,因为它们最大程度地减少了内存使用并提高了性能。

Node.js 中有四种主要的流类型,所有这些都是 EventEmitter 的实例:

  • Readable:可以从中读取数据的流(例如,fs.createReadStream,HTTP 请求)。
  • Writable:可以向其写入数据的流(例如,fs.createWriteStream,HTTP 响应)。
  • Duplex:既是 Readable 又是 Writable 的流(例如,TCP socket,zlib 流)。
  • Transform:一种特殊的 Duplex 流,其输出是根据输入计算得出的(例如,在数据通过时对其进行压缩或加密的流)。

由于流是 EventEmitters,它们在生命周期中会触发各种事件。一些常见的事件包括:

  • 'data':当有数据块可从 Readable 流中读取时触发。
  • 'end':当 Readable 流中没有更多数据可读时触发。
  • 'error':如果在读或写过程中发生错误时触发。
  • 'finish':在 Writable 流上触发,表示所有数据已冲刷(flushed)到底层系统(例如,写入文件)。
  • 'close':当流及其底层资源(如文件描述符)被关闭时触发。

让我们使用 Readable 流读取文件的内容。首先,创建一个文件 input.txt:

This is line one.
This is line two.
Streaming data with Node.js is efficient!

现在,创建 readStream.js:

const fs = require('fs');
let fileContent = '';
// 创建一个可读流
// 直接在 options 中指定编码,以便自动转换为字符串
const readerStream = fs.createReadStream('input.txt', { encoding: 'utf8' });
// 处理流事件:
// 'data' 事件:当数据块到达时触发
readerStream.on('data', (chunk) => {
console.log(`Received ${chunk.length} bytes of data.`);
fileContent += chunk;
});
// 'end' 事件:当没有更多数据可读时触发
readerStream.on('end', () => {
console.log('\n--- End of Stream ---');
console.log('Complete File Content:\n' + fileContent);
});
// 'error' 事件:发生错误时触发
readerStream.on('error', (err) => {
console.error('Error reading stream:', err.stack);
});
console.log('Stream reading initiated...');

运行脚本:

$ node readStream.js

输出(数据块大小可能有所不同):

Stream reading initiated...
Received 73 bytes of data.
--- End of Stream ---
Complete File Content:
This is line one.
This is line two.
Streaming data with Node.js is efficient!

注意:对于非常大的文件,使用 += 拼接数据块可能效率低下。对于大量数据,考虑单独处理数据块或使用专门的流消费者。

让我们使用 Writable 流将数据写入文件。创建 writeStream.js:

const fs = require('fs');
const dataToWrite = 'Writing data chunk by chunk using Writable streams.\nHello Streams!';
// 创建一个可写流指向 output.txt
const writerStream = fs.createWriteStream('output.txt', { encoding: 'utf8' });
// 将数据写入流
writerStream.write(dataToWrite);
// 标记写入结束 - 这对于触发 'finish' 事件至关重要
writerStream.end(); // 也可以传递最后一个数据块:writerStream.end('last chunk');
// 处理流事件:
// 'finish' 事件:当所有数据都被冲刷(flushed)后触发
writerStream.on('finish', () => {
console.log('Write operation completed.');
});
// 'error' 事件:发生错误时触发
writerStream.on('error', (err) => {
console.error('Error writing stream:', err.stack);
});
console.log('Stream writing initiated...');

运行脚本:

$ node writeStream.js

输出:

Stream writing initiated...
Write operation completed.

检查在同一目录中创建的 output.txt 文件。它应该包含 dataToWrite 中定义的内容。

管道连接(Piping)是一种强大的机制,它将 Readable 流的输出直接连接到 Writable 流的输入。它自动处理数据的流动、背压(backpressure,防止 Writable 流被压垮)和错误传播。

.pipe() 方法在 Readable 流上可用。它接受一个 Writable 流作为参数。

示例:使用 pipe 复制文件(pipeDemo.js):

const fs = require('fs');
// 从 input.txt 创建可读流
const readerStream = fs.createReadStream('input.txt');
// 创建可写流指向 output_piped.txt
const writerStream = fs.createWriteStream('output_piped.txt');
// 连接流(Pipe the streams)!
readerStream.pipe(writerStream);
// 可选:监听写入流上的 finish/error 事件
writerStream.on('finish', () => {
console.log('Piping finished: input.txt copied to output_piped.txt');
});
writerStream.on('error', (err) => {
console.error('Error during pipe:', err);
});
readerStream.on('error', (err) => {
console.error('Error reading input for pipe:', err);
});
console.log('Piping initiated...');

运行它:

$ node pipeDemo.js

输出:

Piping initiated...
Piping finished: input.txt copied to output_piped.txt

这会自动高效地将 input.txt 的内容复制到 output_piped.txt。

你可以将多个 .pipe() 调用链接在一起,通常涉及 Transform 流,以创建数据处理管道。

示例:使用 zlib(一个 Transform 流)压缩文件(compress.js):

const fs = require('fs');
const zlib = require('zlib'); // 内置的压缩模块
const { pipeline } = require('stream'); // 推荐用于流连接(piping)的方式
// 创建流
const reader = fs.createReadStream('input.txt');
const gzip = zlib.createGzip(); // Gzip 压缩转换流
const writer = fs.createWriteStream('input.txt.gz');
console.log('Starting file compression...');
// 使用 stream.pipeline 实现健壮的错误处理和清理
pipeline(reader, gzip, writer, (err) => {
if (err) {
console.error('Pipeline failed:', err);
} else {
console.log('File compressed successfully to input.txt.gz');
}
});

运行它:

$ node compress.js

输出:

Starting file compression...
File compressed successfully to input.txt.gz

这将创建一个压缩文件 input.txt.gz。

示例:解压缩文件(decompress.js):

const fs = require('fs');
const zlib = require('zlib');
const { pipeline } = require('stream');
// 创建流
const reader = fs.createReadStream('input.txt.gz');
const gunzip = zlib.createGunzip(); // Gunzip 解压缩转换流
const writer = fs.createWriteStream('input_decompressed.txt');
console.log('Starting file decompression...');
// 使用 stream.pipeline 实现健壮的错误处理和清理
pipeline(reader, gunzip, writer, (err) => {
if (err) {
console.error('Pipeline failed:', err);
} else {
console.log('File decompressed successfully to input_decompressed.txt');
}
});

运行它:

$ node decompress.js

输出:

Starting file decompression...
File decompressed successfully to input_decompressed.txt

这将 input.txt.gz 解压缩回文本文件。

使用 stream.pipeline 通常比多个 .pipe() 调用更受欢迎,因为它提供了更好的错误处理,并且在管道中的某个流失败时会自动销毁其他流。

流是一个深入的主题。请查阅 Node.js 官方文档以获取更多详细信息: