引言
Stream 是 Node.js 最核心的抽象之一。从文件读写、HTTP 请求响应到数据库连接,Stream 无处不在。然而,许多开发者对 Stream 的理解停留在"pipe 一下就好"的层面,对背压(backpressure)、错误传播、自定义流的实现知之甚少。
本文将深入 Stream 的内部机制,从四种基本流类型出发,结合实际场景,帮你真正掌握这一强大工具。
一、为什么需要 Stream?
传统方式的局限
假设你需要读取一个 2GB 的日志文件并写入到另一个位置:
javascript
// ❌ 传统方式:将整个文件读入内存
const fs = require('fs');
fs.readFile('/path/to/huge-file.log', (err, data) => {
if (err) throw err;
// 2GB 数据全部在内存中 —— 大多数服务器根本无法承受
fs.writeFile('/path/to/output.log', data, (err) => {
if (err) throw err;
});
});Stream 的解决方案
javascript
// ✅ Stream 方式:分块处理,内存占用恒定
const fs = require('fs');
const readStream = fs.createReadStream('/path/to/huge-file.log', {
highWaterMark: 64 * 1024, // 每次读取 64KB
});
const writeStream = fs.createWriteStream('/path/to/output.log');
readStream.pipe(writeStream);
readStream.on('error', (err) => console.error('读取失败:', err));
writeStream.on('error', (err) => console.error('写入失败:', err));
writeStream.on('finish', () => console.log('文件复制完成'));二、四种基本流类型
Node.js 提供了四种流抽象,覆盖了所有数据流动场景:
| 类型 | 方向 | 核心方法 | 典型场景 |
|---|---|---|---|
| Readable | 数据流出 | _read() |
文件读取、HTTP 请求体 |
| Writable | 数据流入 | _write() |
文件写入、HTTP 响应 |
| Duplex | 双向(独立) | _read() + _write() |
TCP Socket、WebSocket |
| Transform | 双向(转换) | _transform() |
压缩、加密、解析 |
Readable Stream
Readable 流负责生产数据。数据以两种模式流动:
javascript
const { Readable } = require('stream');
// 实现一个生成斐波那契数列的 Readable 流
class FibonacciStream extends Readable {
constructor(max = 1000) {
super({ objectMode: true });
this.max = max;
this.a = 0;
this.b = 1;
}
_read() {
const next = this.a + this.b;
this.a = this.b;
this.b = next;
if (next > this.max) {
this.push(null); // 结束流
return;
}
this.push({ value: next, index: this.index++ });
}
}
const fib = new FibonacciStream(10000);
fib.on('data', (chunk) => {
console.log(chunk);
});
fib.on('end', () => {
console.log('数列生成完毕');
});Writable Stream
javascript
const { Writable } = require('stream');
// 实现一个批量写入数据库的 Writable 流
class BatchInsertStream extends Writable {
constructor(db, options) {
super({ objectMode: true, ...options });
this.db = db;
this.buffer = [];
this.batchSize = 100;
}
async _write(chunk, encoding, callback) {
this.buffer.push(chunk);
if (this.buffer.length >= this.batchSize) {
try {
await this.db.insertMany(this.buffer);
this.buffer = [];
callback(); // 写入成功
} catch (err) {
callback(err); // 传播错误
}
} else {
callback();
}
}
async _final(callback) {
// 处理剩余的缓冲区数据
try {
if (this.buffer.length > 0) {
await this.db.insertMany(this.buffer);
}
callback();
} catch (err) {
callback(err);
}
}
}Transform Stream
Transform 流是最常用也最强大的流类型,同时实现了可读和可写接口:
javascript
const { Transform } = require('stream');
// JSON 行解析 Transform 流
class JSONLineParser extends Transform {
constructor(options) {
super({ readableObjectMode: true, ...options });
this.buffer = '';
}
_transform(chunk, encoding, callback) {
this.buffer += chunk.toString();
// 按行分割并解析 JSON
const lines = this.buffer.split('\n');
// 最后一行可能不完整,保留在下一次处理
this.buffer = lines.pop() || '';
for (const line of lines) {
if (line.trim()) {
try {
this.push(JSON.parse(line));
} catch (err) {
callback(err);
return;
}
}
}
callback();
}
_flush(callback) {
// 处理最终的缓冲区
if (this.buffer.trim()) {
try {
this.push(JSON.parse(this.buffer));
} catch (err) {
callback(err);
return;
}
}
callback();
}
}三、背压机制深度解析
背压(backpressure)是 Stream 最容易被忽视但最关键的概念。当数据生产速度超过消费速度时,如果放任不管,内存会被无限增长的内部缓冲区撑爆。
背压的工作原理
javascript
const fs = require('fs');
const readable = fs.createReadStream('source.dat', {
highWaterMark: 16 * 1024, // 16KB 内部缓冲上限
});
const writable = fs.createWriteStream('dest.dat');
// readable.pipe(writable) 内部的实际流程:
readable.on('data', (chunk) => {
// writable.write() 返回 false 表示内部缓冲区已满
const canContinue = writable.write(chunk);
if (!canContinue) {
// 暂停读取,等待 writable 排空
readable.pause();
// 当 writable 排空时恢复读取
writable.once('drain', () => {
readable.resume();
});
}
});
writable.on('drain', () => {
console.log('缓冲区排空,ready 继续写入');
});pipe() 的完整手动实现
理解 pipe() 的实现有助于掌握背压:
javascript
function manualPipe(readable, writable) {
// 错误处理:任一流出错,销毁另一方
const onReadableError = (err) => {
writable.destroy(err);
};
const onWritableError = (err) => {
readable.destroy(err);
};
readable.on('error', onReadableError);
writable.on('error', onWritableError);
readable.on('data', (chunk) => {
const canWrite = writable.write(chunk);
if (!canWrite) {
readable.pause();
// writable 排空后恢复 readable
writable.once('drain', () => {
readable.resume();
});
}
});
// readable 结束时关闭 writable
readable.on('end', () => {
writable.end();
});
return writable;
}highWaterMark 调优
javascript
// 不同场景下的 highWaterMark 选择
const fs = require('fs');
// 大文件复制:较大的缓冲区减少 I/O 调用次数
const bigFileRead = fs.createReadStream('large-video.mp4', {
highWaterMark: 1024 * 1024, // 1MB
});
// 实时日志处理:较小的缓冲区降低延迟
const logRead = fs.createReadStream('app.log', {
highWaterMark: 4 * 1024, // 4KB
});
// 对象模式:highWaterMark 是对象数量而非字节
const { Readable } = require('stream');
const objStream = new Readable({
objectMode: true,
highWaterMark: 16, // 16 个对象
read() {},
});四、pipeline:现代 Stream 编排
stream.pipeline 是 pipe() 的现代替代品,提供了统一的错误处理:
javascript
const { pipeline } = require('stream/promises');
const fs = require('fs');
const zlib = require('zlib');
const crypto = require('crypto');
async function compressAndEncrypt(inputPath, outputPath, password) {
try {
await pipeline(
fs.createReadStream(inputPath),
zlib.createGzip(),
crypto.createCipheriv('aes-256-gcm', key, iv),
fs.createWriteStream(outputPath),
);
console.log('压缩加密完成');
} catch (err) {
// pipeline 会自动销毁所有流并传播错误
console.error('管道处理失败:', err);
}
}pipeline vs pipe 对比
| 特性 | pipe() | pipeline() |
|---|---|---|
| 错误传播 | 不自动传播 | 自动传播到回调/Promise |
| 流清理 | 需手动 destroy | 自动销毁所有中间流 |
| 回调/Promise | 仅 EventEmitter | 支持 callback 和 Promise |
| 可组合性 | 链式调用 | 数组形式,可读性更好 |
五、实战案例
案例 1:CSV 大文件解析与数据库导入
javascript
const fs = require('fs');
const { pipeline } = require('stream/promises');
const { Transform } = require('stream');
// CSV 行解析器
class CSVParser extends Transform {
constructor(headers, options) {
super({ readableObjectMode: true, ...options });
this.headers = headers;
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 values = line.split(',');
const record = {};
this.headers.forEach((header, i) => {
record[header] = values[i]?.trim() ?? '';
});
this.push(record);
}
callback();
}
_flush(callback) {
if (this.buffer.trim()) {
const values = this.buffer.split(',');
const record = {};
this.headers.forEach((header, i) => {
record[header] = values[i]?.trim() ?? '';
});
this.push(record);
}
callback();
}
}
// 批量插入 Transform
class BatchInsertTransformer extends Transform {
constructor(db, batchSize = 1000) {
super({ objectMode: true });
this.db = db;
this.batchSize = batchSize;
this.batch = [];
}
async _transform(record, encoding, callback) {
this.batch.push(record);
if (this.batch.length >= this.batchSize) {
try {
await this.db.insertMany(this.batch);
console.log(`已插入 ${this.batch.length} 条记录`);
this.batch = [];
callback();
} catch (err) {
callback(err);
}
} else {
callback();
}
}
async _flush(callback) {
try {
if (this.batch.length > 0) {
await this.db.insertMany(this.batch);
console.log(`已插入最后 ${this.batch.length} 条记录`);
}
callback();
} catch (err) {
callback(err);
}
}
}
// 使用
async function importCSV(filePath, db) {
const headers = ['id', 'name', 'email', 'created_at', 'status'];
await pipeline(
fs.createReadStream(filePath, { highWaterMark: 64 * 1024 }),
new CSVParser(headers),
new BatchInsertTransformer(db, 500),
);
console.log('导入完成');
}案例 2:实时日志过滤与转发
javascript
const { pipeline } = require('stream/promises');
const { Transform } = require('stream');
const net = require('net');
// 日志级别过滤器
class LogFilter extends Transform {
constructor(minLevel, options) {
super({ ...options });
this.minLevel = minLevel;
}
_transform(chunk, encoding, callback) {
const line = chunk.toString().trim();
// 检查日志级别
const levelMatch = line.match(/\[(ERROR|WARN|INFO|DEBUG)\]/);
if (levelMatch) {
const level = levelMatch[1];
const shouldForward = this.shouldLog(level);
if (shouldForward) {
this.push(chunk);
}
}
callback();
}
shouldLog(level) {
const levels = { DEBUG: 0, INFO: 1, WARN: 2, ERROR: 3 };
return levels[level] >= levels[this.minLevel];
}
}
// 使用标准输入作为日志源
async function filterStdin() {
await pipeline(
process.stdin,
new LogFilter('WARN'),
process.stdout,
);
}
// 命令行使用:
// tail -f /var/log/app.log | node log-filter.js案例 3:HTTP 大文件下载(支持断点续传)
javascript
const http = require('http');
const fs = require('fs');
const { pipeline } = require('stream/promises');
async function downloadFile(url, destPath) {
return new Promise((resolve, reject) => {
http.get(url, (response) => {
// 处理重定向
if (response.statusCode >= 300 && response.statusCode < 400) {
return downloadFile(response.headers.location, destPath)
.then(resolve)
.catch(reject);
}
if (response.statusCode !== 200) {
reject(new Error(`下载失败: HTTP ${response.statusCode}`));
return;
}
const totalSize = parseInt(response.headers['content-length'], 10);
let downloadedSize = 0;
// 使用 Transform 流监控进度
const progressTracker = new Transform({
transform(chunk, encoding, callback) {
downloadedSize += chunk.length;
const progress = ((downloadedSize / totalSize) * 100).toFixed(1);
process.stdout.write(`\r下载进度: ${progress}%`);
this.push(chunk);
callback();
},
});
pipeline(
response,
progressTracker,
fs.createWriteStream(destPath),
)
.then(() => {
console.log('\n下载完成');
resolve();
})
.catch(reject);
}).on('error', reject);
});
}六、性能与内存优化
内存使用对比
javascript
const fs = require('fs');
// 测试:处理 1GB 文件的不同方式内存占用
// 方式 A:readFile —— 内存 ≈ 1GB + overhead
// fs.readFile('1gb-file.dat', callback);
// 方式 B:Stream —— 内存 ≈ 64KB(highWaterMark)
// fs.createReadStream('1gb-file.dat', { highWaterMark: 64 * 1024 });
// 方式 C:Stream + 处理 —— 内存 ≈ 64KB + 处理开销
// stream.pipe(transform).pipe(writeStream);并发流控制
javascript
const { pipeline } = require('stream/promises');
const { Transform } = require('stream');
// 并发处理 Transform:允许 N 个异步任务同时进行
class ParallelTransform extends Transform {
constructor(concurrency, transformFn, options) {
super({ objectMode: true, ...options });
this.concurrency = concurrency;
this.transformFn = transformFn;
this.active = 0;
this.queue = [];
}
_transform(chunk, encoding, callback) {
this.active++;
this.transformFn(chunk)
.then((result) => {
this.push(result);
this.active--;
this._processQueue();
})
.catch(callback);
if (this.active < this.concurrency) {
callback();
} else {
this.queue.push(callback);
}
}
_processQueue() {
while (this.active < this.concurrency && this.queue.length > 0) {
const cb = this.queue.shift();
cb();
}
}
}
// 使用:并发 5 个请求批量处理记录
async function processWithConcurrency() {
const transform = new ParallelTransform(5, async (record) => {
const result = await fetch(`https://api.example.com/process`, {
method: 'POST',
body: JSON.stringify(record),
});
return result.json();
});
await pipeline(
sourceStream,
transform,
destinationStream,
);
}七、常见陷阱与最佳实践
1. 忘记处理 error 事件
javascript
// ❌ 未处理的 error 会导致进程崩溃
readStream.pipe(writeStream);
// ✅ 始终处理错误
readStream.on('error', handleError);
writeStream.on('error', handleError);
// ✅ 使用 pipeline 自动传播错误
await pipeline(readStream, writeStream);2. 混合使用 pipe 和事件
javascript
// ❌ 同时使用 pipe 和 data 事件
readStream.pipe(writeStream);
readStream.on('data', (chunk) => {
// 数据会被两个消费者竞争
});
// ✅ 需要观测数据时使用 Transform 流
const observer = new Transform({
transform(chunk, encoding, callback) {
console.log(`处理了 ${chunk.length} 字节`);
this.push(chunk);
callback();
},
});
readStream.pipe(observer).pipe(writeStream);3. 流完成后的清理
javascript
async function processFile(input, output) {
const readStream = fs.createReadStream(input);
const writeStream = fs.createWriteStream(output);
try {
await pipeline(readStream, writeStream);
// 流已自动关闭,无需额外清理
} catch (err) {
// pipeline 失败时自动销毁所有流
console.error('处理失败:', err);
}
}结语
Stream 是 Node.js 的基石——它不仅是 API,更是一种编程范式。掌握 Stream 意味着你可以在处理数据时保持恒定的内存占用、天然的背压支持和优雅的组合能力。
核心要点:
- 四种流类型对应不同的数据流向需求,理解
_read、_write、_transform是关键 - 背压是 Stream 设计中最精妙的部分,
highWaterMark和drain事件是核心机制 - pipeline 取代
pipe()成为现代准则,提供自动错误传播和流清理 - 自定义流让你可以像搭积木一样组合数据处理管道
- 始终处理错误——未捕获的 Stream error 会直接崩溃 Node.js 进程
当你下次面对大文件、实时数据或高并发场景时,首先问自己:这能不能用 Stream 解决?