1. 为什么需要 Stream
Stream 是 Node.js 的核心抽象,处理大量数据时无需一次性加载到内存,边读边处理,内存占用恒定。
无 Stream(全部加载):
读取 1GB 文件 → 加载到内存 → 处理 → 写入新文件
内存占用: 1GB+
有 Stream(边读边写):
读取 chunk ──→ 处理 chunk ──→ 写入 chunk
内存占用: 16KB (chunk size)
2. 四种流类型
| 类型 | 说明 | 示例 |
|---|---|---|
| Readable | 可读流(数据源) | fs.createReadStream、HTTP req |
| Writable | 可写流(数据目标) | fs.createWriteStream、HTTP res |
| Duplex | 可读可写(双向) | net.Socket、TCP 连接 |
| Transform | 转换流(读入→转换→写出) | zlib.createGzip、加密流 |
2.1 Readable 流
const fs = require('fs');
// 创建可读流
const readable = fs.createReadStream('large-file.txt', {
encoding: 'utf8',
highWaterMark: 16 * 1024 // 每次读取 16KB
});
// 消费数据
readable.on('data', (chunk) => {
console.log(`Received ${chunk.length} bytes`);
});
readable.on('end', () => console.log('Done'));
readable.on('error', (err) => console.error(err));
// 或使用异步迭代器(更现代)
(async () => {
for await (const chunk of readable) {
console.log(chunk);
}
})();
2.2 Writable 流
const writable = fs.createWriteStream('output.txt');
writable.write('Hello ');
writable.write('World\n');
writable.end('Done!');
writable.on('finish', () => console.log('写入完成'));
2.3 Transform 转换流
const { Transform } = require('stream');
// 自定义转换流:将所有输入转为大写
const upperCaseTransform = new Transform({
transform(chunk, encoding, callback) {
this.push(chunk.toString().toUpperCase());
callback();
}
});
// 使用管道链
fs.createReadStream('input.txt')
.pipe(upperCaseTransform)
.pipe(fs.createWriteStream('output.txt'));
3. 管道与背压
3.1 Pipe 管道链
const fs = require('fs');
const zlib = require('zlib');
// 压缩文件流式处理
fs.createReadStream('input.txt')
.pipe(zlib.createGzip()) // 压缩
.pipe(fs.createWriteStream('input.txt.gz'));
// 解压
fs.createReadStream('input.txt.gz')
.pipe(zlib.createGunzip())
.pipe(fs.createWriteStream('output.txt'));
3.2 背压(Backpressure)
当 Writable 消费速度慢于 Readable 生产速度时,需要背压机制防止内存无限增长。
const readable = fs.createReadStream('big-file.txt');
const writable = fs.createWriteStream('output.txt');
// pipe 自动处理背压:
// readable 检测到 writable.write() 返回 false → 暂停读取
// writable 的 drain 事件触发 → 恢复读取
readable.pipe(writable);
// 手动处理背压
readable.on('data', (chunk) => {
const canContinue = writable.write(chunk);
if (!canContinue) {
readable.pause(); // 暂停读取
writable.once('drain', () => readable.resume()); // 等 writable 就绪
}
});
4. Buffer 二进制处理
Buffer 是 Node.js 处理二进制数据的类,固定大小、直接操作内存。
// 创建 Buffer
const buf1 = Buffer.alloc(10); // 10 字节,初始化为 0
const buf2 = Buffer.from('Hello'); // 从字符串创建
const buf3 = Buffer.from([0x48, 0x65]); // 从数组创建
// 读写
buf1.write('Node.js');
console.log(buf1.toString()); // 'Node.js'
console.log(buf1.toString('hex')); // 十六进制
console.log(buf1.toString('base64')); // Base64
// 切片(共享内存,不复制)
const slice = buf2.slice(0, 3);
slice[0] = 0x41; // 修改 slice 会影响原 buf2!
// 拼接
const buf4 = Buffer.concat([buf2, buf3]);
// 编码转换
const utf8 = Buffer.from('你好', 'utf8');
console.log(utf8.toString('utf8')); // '你好'
延伸阅读
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。