引言
在 Node.js 开发中,处理大文件或大量数据时,传统的一次性读取方式(如 fs.readFile)往往会导致内存溢出或性能瓶颈。流(Stream)是 Node.js 中处理数据的核心抽象,它允许你以分块的方式处理数据,从而节省内存并提高效率。本文将带你从零开始,掌握 Node.js Stream 的实战技巧,包括读写流、转换流、管道操作以及背压处理。
1. Stream 基础
Node.js 的流有四种类型:
- Readable:可读流(如
fs.createReadStream) - Writable:可写流(如
fs.createWriteStream) - Transform:转换流(如
zlib.createGzip) - Duplex:双工流(可读可写,如
net.Socket)
所有流都是 EventEmitter 的实例,通过事件驱动。
1.1 创建简单的读写流
const fs = require('fs');
// 创建可读流
const readableStream = fs.createReadStream('input.txt', {
highWaterMark: 16 * 1024, // 每次读取 16KB,默认 64KB
encoding: 'utf8'
});
// 创建可写流
const writableStream = fs.createWriteStream('output.txt');
// 监听数据事件
readableStream.on('data', (chunk) => {
console.log(`Received ${chunk.length} bytes of data.`);
writableStream.write(chunk);
});
readableStream.on('end', () => {
console.log('Reading finished.');
writableStream.end();
});
// 错误处理
readableStream.on('error', (err) => console.error('Read error:', err));
writableStream.on('error', (err) => console.error('Write error:', err));
注意:手动管理读写容易出错,比如忘记调用 writableStream.end() 会导致文件不完整。推荐使用 pipe 方法。
2. 管道(Pipe)操作
pipe 是流最强大的功能,它自动管理背压(backpressure),防止数据积压导致内存溢出。
const fs = require('fs');
const readable = fs.createReadStream('largefile.txt');
const writable = fs.createWriteStream('copy.txt');
readable.pipe(writable);
writable.on('finish', () => console.log('File copied successfully.'));
2.1 链式管道
管道可以链式调用,形成数据处理流水线。
const fs = require('fs');
const zlib = require('zlib');
// 压缩文件
fs.createReadStream('input.txt')
.pipe(zlib.createGzip())
.pipe(fs.createWriteStream('input.txt.gz'))
.on('finish', () => console.log('Compression done.'));
💡 最佳实践:始终在最后一个流上监听 finish 或 close 事件,确保所有数据已处理完毕。
3. 自定义转换流
转换流用于在数据传递过程中修改数据。例如,创建一个将文本转为大写的流。
const { Transform } = require('stream');
class UpperCaseTransform extends Transform {
_transform(chunk, encoding, callback) {
// chunk 是 Buffer 或字符串,取决于 encoding
this.push(chunk.toString().toUpperCase());
callback();
}
}
// 使用
const upperCase = new UpperCaseTransform();
process.stdin.pipe(upperCase).pipe(process.stdout);
3.1 实战:CSV 数据清洗
假设我们需要读取一个大型 CSV 文件,过滤掉某些行,并输出为 JSON。
const fs = require('fs');
const { Transform } = require('stream');
// 创建一个过滤转换流
class FilterTransform extends Transform {
constructor(options) {
super({ objectMode: true, ...options }); // 使用对象模式
this.header = null;
}
_transform(chunk, encoding, callback) {
// 假设 chunk 是一行字符串
const line = chunk.toString().trim();
if (!line) return callback();
if (!this.header) {
this.header = line.split(',');
return callback();
}
const values = line.split(',');
if (values[1] === 'valid') { // 假设第二列是状态
const obj = {};
this.header.forEach((key, i) => { obj[key] = values[i]; });
this.push(JSON.stringify(obj) + '\n');
}
callback();
}
}
// 使用
const readStream = fs.createReadStream('data.csv');
const writeStream = fs.createWriteStream('output.json');
readStream
.pipe(new FilterTransform())
.pipe(writeStream)
.on('finish', () => console.log('CSV processed.'));
注意:使用 objectMode 后,流处理的是对象而非 Buffer,但要注意内存管理。
4. 背压(Backpressure)详解
背压是指当写入速度慢于读取速度时,流自动控制数据流动的机制。pipe 方法自动处理背压,但手动实现时需注意。
4.1 手动处理背压
const fs = require('fs');
const readable = fs.createReadStream('large.txt');
const writable = fs.createWriteStream('copy.txt');
readable.on('data', (chunk) => {
const canContinue = writable.write(chunk);
if (!canContinue) {
console.log('Backpressure detected, pausing read.');
readable.pause();
writable.once('drain', () => {
console.log('Drain event, resuming read.');
readable.resume();
});
}
});
readable.on('end', () => writable.end());
💡 最佳实践:除非有特殊需求,否则始终使用 pipe 或 pipeline(Node.js 10+)来处理背压。
5. 使用 pipeline 替代 pipe
pipeline 是 Node.js 10 引入的模块化方法,它自动处理错误和清理。
const { pipeline } = require('stream');
const fs = require('fs');
const zlib = require('zlib');
pipeline(
fs.createReadStream('input.txt'),
zlib.createGzip(),
fs.createWriteStream('input.txt.gz'),
(err) => {
if (err) {
console.error('Pipeline failed:', err);
} else {
console.log('Pipeline succeeded.');
}
}
);
pipeline 在发生错误时会自动销毁所有流,避免资源泄漏。
6. 实战:大文件行处理
处理大文件时,逐行读取比按 chunk 读取更方便。但 Node.js 没有内置的行读取流,我们可以自己实现或使用 readline 模块。
const fs = require('fs');
const readline = require('readline');
async function processLargeFile(filePath) {
const fileStream = fs.createReadStream(filePath);
const rl = readline.createInterface({
input: fileStream,
crlfDelay: Infinity
});
let lineCount = 0;
for await (const line of rl) {
// 处理每一行
lineCount++;
if (lineCount % 100000 === 0) {
console.log(`Processed ${lineCount} lines`);
}
}
console.log(`Total lines: ${lineCount}`);
}
processLargeFile('huge.log').catch(console.error);
注意:readline 内部使用了流,但它是按行分割的,适合处理文本文件。
7. 常见陷阱与最佳实践
7.1 不要忘记错误处理
流事件中的错误必须被捕获,否则会导致进程崩溃。
const fs = require('fs');
const stream = fs.createReadStream('nonexistent.txt');
stream.on('error', (err) => console.error('Error:', err.message));
7.2 避免在 ‘data’ 事件中执行同步操作
同步操作会阻塞事件循环,影响性能。
7.3 使用 highWaterMark 控制缓冲区大小
根据内存和性能需求调整 highWaterMark。
7.4 流不是万能的
对于小文件(< 50MB),直接使用 fs.readFile 可能更简单。
总结
Node.js Stream 是处理大文件和构建数据管道的利器。通过 pipe 或 pipeline,你可以轻松组合多个流,实现高效的数据处理。本文介绍了流的基本用法、自定义转换流、背压处理以及实战案例。掌握这些技巧后,你将能够优雅地处理 GB 级别的数据。
延伸阅读:
- Node.js 官方文档:Stream API
- 学习使用
stream.finished()和stream.destroy()进行资源清理 - 探索第三方流库如
through2、split2等