为什么你需要掌握 Stream?
在 Node.js 中处理大文件(如几 GB 的日志、视频或数据集)时,传统的一次性读取所有内容到内存的方式会轻易耗尽内存,甚至导致进程崩溃。Stream(流)正是为此而生——它将数据分块处理,每次只占用一小块内存,从而高效处理任意大小的数据。
然而,很多开发者对 Stream 的认知停留在“用 pipe 连接”的层面,遇到背压、错误处理、自定义流等问题时常常手足无措。本文将从实战出发,带你深入理解 Stream 的核心机制,并给出可复用的最佳实践。
Stream 基础回顾
Node.js 中有四种基本的流类型:
- Readable(可读流):数据源,如
fs.createReadStream - Writable(可写流):数据目标,如
fs.createWriteStream - Transform(转换流):既可读又可写,可在读写过程中修改数据,如
zlib.createGzip - Duplex(双工流):既可读又可写,但读写独立,如
net.Socket
核心方法:
readable.pipe(writable):将可读流连接到可写流,自动处理背压readable.on('data', chunk => ...):手动消费数据(不推荐,除非需要精细控制)writable.write(chunk):写入数据块writable.end():结束写入
实战一:大文件复制与背压控制
错误示例:一次性读取
const fs = require('fs');
// 危险!大文件会耗尽内存
fs.readFile('bigfile.txt', (err, data) => {
fs.writeFile('copy.txt', data, err => {
if (err) console.error(err);
});
});
正确做法:使用 pipe 自动背压
const fs = require('fs');
const readable = fs.createReadStream('bigfile.txt');
const writable = fs.createWriteStream('copy.txt');
readable.pipe(writable);
readable.on('error', err => console.error('读取错误:', err));
writable.on('error', err => console.error('写入错误:', err));
writable.on('finish', () => console.log('复制完成'));
背压(Backpressure):当写入速度慢于读取速度时,pipe 会自动暂停读取,等待写入完成后再继续。这是 Stream 最强大的特性之一。
手动控制背压(高级场景)
有时你需要更精细的控制,比如统计进度或转换数据。此时可以监听 drain 事件:
const fs = require('fs');
const readable = fs.createReadStream('bigfile.txt');
const writable = fs.createWriteStream('copy.txt');
let totalBytes = 0;
readable.on('data', chunk => {
totalBytes += chunk.length;
// 如果 write 返回 false,表示缓冲区已满,应暂停读取
if (!writable.write(chunk)) {
readable.pause();
}
console.log(`已处理 ${totalBytes} 字节`);
});
writable.on('drain', () => {
readable.resume();
});
readable.on('end', () => writable.end());
readable.on('error', err => console.error(err));
writable.on('finish', () => console.log('复制完成'));
> 💡 注意:手动实现背压时,务必确保在 drain 事件中恢复读取,否则可能导致内存泄漏或程序挂起。
实战二:自定义 Transform 流实现数据转换
假设我们需要将 CSV 文件中的每一行转成 JSON 对象,并过滤掉某些字段。我们可以创建一个 Transform 流:
const { Transform } = require('stream');
const fs = require('fs');
class CsvToJsonTransform extends Transform {
constructor(options = {}) {
super(options);
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() === '') continue;
const [name, age, city] = line.split(',');
const json = JSON.stringify({ name, age: parseInt(age), city }) + '\n';
this.push(json);
}
callback();
}
_flush(callback) {
// 处理最后剩余的行
if (this.buffer.trim() !== '') {
const [name, age, city] = this.buffer.split(',');
const json = JSON.stringify({ name, age: parseInt(age), city }) + '\n';
this.push(json);
}
callback();
}
}
// 使用
const readable = fs.createReadStream('data.csv');
const writable = fs.createWriteStream('data.json');
const transformer = new CsvToJsonTransform();
readable.pipe(transformer).pipe(writable);
readable.on('error', err => console.error(err));
transformer.on('error', err => console.error(err));
writable.on('finish', () => console.log('转换完成'));
> 💡 关键点: > – transform 方法处理每个数据块,注意处理跨块的行(buffer 机制) > – flush 方法在流结束时被调用,用于处理剩余数据 > – 使用 this.push() 输出转换后的数据
实战三:数据管道与错误处理最佳实践
实际项目中,我们经常需要串联多个流(如解压、解密、转换、压缩)。错误处理是容易忽略的环节,一个流的错误如果不处理,会传播到整个管道,导致程序崩溃。
管道中的错误传播
const fs = require('fs');
const zlib = require('zlib');
const { pipeline } = require('stream');
const readable = fs.createReadStream('archive.gz');
const gunzip = zlib.createGunzip();
const writable = fs.createWriteStream('output.txt');
// 使用 pipeline 而不是 pipe,自动处理错误和清理
pipeline(
readable,
gunzip,
writable,
(err) => {
if (err) {
console.error('Pipeline 失败:', err);
} else {
console.log('Pipeline 成功');
}
}
);
为什么推荐 pipeline?
pipe不会自动销毁流,错误可能导致内存泄漏pipeline在 Node.js 10+ 中引入,会自动销毁所有流并调用回调- 手动使用
pipe需要监听每个流的错误并手动销毁,非常繁琐
自定义错误处理与重试
对于网络请求等场景,我们可以封装一个可重试的流:
const { Transform } = require('stream');
class RetryTransform extends Transform {
constructor(options = {}) {
super(options);
this.retries = options.retries || 3;
}
_transform(chunk, encoding, callback) {
// 模拟一个可能失败的操作
this.processWithRetry(chunk, 0, callback);
}
processWithRetry(chunk, attempt, callback) {
try {
// 假设这里可能抛出错误
if (Math.random() < 0.3) throw new Error('随机失败');
const result = chunk.toString().toUpperCase();
this.push(result);
callback();
} catch (err) {
if (attempt < this.retries) {
console.log(`重试第 ${attempt + 1} 次`);
setTimeout(() => this.processWithRetry(chunk, attempt + 1, callback), 1000);
} else {
callback(err); // 超过重试次数,传递错误
}
}
}
}
> ⚠️ 注意:重试逻辑应谨慎使用,避免无限重试导致资源耗尽。
实战四:对象模式流(Object Mode)
默认情况下,流处理的是 Buffer 或字符串。但有时我们需要处理 JavaScript 对象,这时需要启用对象模式:
const { Transform } = require('stream');
const objectStream = new Transform({
objectMode: true,
transform(chunk, encoding, callback) {
// chunk 是一个对象
chunk.timestamp = Date.now();
this.push(chunk);
callback();
}
});
objectStream.write({ name: 'Alice' });
objectStream.on('data', obj => console.log(obj));
// 输出: { name: 'Alice', timestamp: 1700000000000 }
> 💡 注意:对象模式会大幅降低性能,因为每个对象都是一个独立的数据块,失去了 Buffer 的批处理优势。仅当需要处理结构化数据且数据量不大时使用。
性能优化与监控
高水位线(HighWaterMark)
highWaterMark 控制流内部缓冲区的最大字节数(或对象数)。调整它可以影响背压行为:
// 增大缓冲区以减少背压触发频率(适合高速读取)
const readable = fs.createReadStream('bigfile.txt', { highWaterMark: 1024 * 1024 }); // 1MB
// 减小缓冲区以降低内存占用(适合低速写入)
const writable = fs.createWriteStream('output.txt', { highWaterMark: 16 * 1024 }); // 16KB
使用 perf_hooks 监控流性能
const { performance, PerformanceObserver } = require('perf_hooks');
const fs = require('fs');
const obs = new PerformanceObserver((items) => {
console.log(items.getEntries()[0]);
});
obs.observe({ entryTypes: ['measure'] });
performance.mark('start');
const readable = fs.createReadStream('bigfile.txt');
const writable = fs.createWriteStream('copy.txt');
readable.pipe(writable);
writable.on('finish', () => {
performance.mark('end');
performance.measure('copy', 'start', 'end');
});
总结
- 永远使用 Stream 处理大文件,避免
readFile/writeFile - 优先使用
pipeline而不是pipe,自动处理错误和清理 - 自定义 Transform 流时,注意跨块数据处理和
_flush - 理解背压机制,手动控制时务必监听
drain事件 - 对象模式只在小数据量时使用,注意性能开销
- 调整
highWaterMark可以优化特定场景的性能
延伸阅读
- Node.js 官方文档:Stream
- 使用
stream.finished检测流是否结束 - 考虑使用
pump库(pipeline 的前身)兼容旧版本 - 探索
readable.pipe(destination, { end: false })实现多路复用