Node\.js Stream 核心:Readable/Writable/Transform \+ 背压问题
Node.js Stream 核心:Readable/Writable/Transform + 背压问题
Node.js Stream(流)是处理大量数据的核心 API,核心优势是边读边处理、不一次性加载全部数据到内存,避免内存溢出。
先明确三大核心流的定位,再讲背压问题和解决方案,这是 Node.js 流的高频面试 + 实战重点。
一、三大核心流的使用场景
Node.js 流基于 stream 模块,最常用的是 Readable(可读流)、Writable(可写流)、Transform(转换流),还有 Duplex(双工流,读写独立,很少直接用)。
1. Readable 可读流
定义:数据源,负责产生 / 读取数据,向外流出数据。
- 核心:你可以从它里面读取数据。
使用场景:
-
读取大文件(日志、视频、CSV),不占用大量内存
-
网络请求接收数据(HTTP 请求体、WebSocket 消息)
-
读取数据库 / 接口返回的流式数据
-
自定义数据源(生成随机数、实时日志推送)
简单示例:
const { createReadStream } = require('fs');
// 读取大文件,分块读取
const readStream = createReadStream('large-file.txt', { encoding: 'utf8' });
readStream.on('data', (chunk) => {
console.log('读取到数据块:', chunk); // 每次读取一小块数据
});
readStream.on('end', () => console.log('读取完成'));
2. Writable 可写流
定义:数据目的地,负责接收 / 写入数据,向内流入数据。
- 核心:你可以向它写入数据。
使用场景:
-
写入大文件(日志落盘、导出大数据)
-
网络响应发送数据(HTTP 响应、文件下载)
-
写入数据库 / 消息队列(流式写入)
简单示例:
const { createWriteStream } = require('fs');
const writeStream = createWriteStream('output.txt');
writeStream.write('第一行数据\n');
writeStream.write('第二行数据\n');
writeStream.end('最后一行数据'); // 结束写入
3. Transform 转换流
定义:中间处理层,可读 + 可写,数据流经它时被修改 / 处理,再输出。
- 核心:读 → 处理 → 写,是数据的 “加工管道”。
使用场景:
-
数据压缩 / 解压(gzip)
-
数据加密 / 解密
-
文本替换、格式转换(CSV 转 JSON、大写转小写)
-
日志过滤、数据清洗
简单示例:
const { Transform } = require('stream');
// 自定义转换流:把文字转大写
const upperTransform = new Transform({
transform(chunk, encoding, callback) {
this.push(chunk.toString().toUpperCase()); // 处理后推送数据
callback();
}
});
// 管道:读 → 转大写 → 写
process.stdin.pipe(upperTransform).pipe(process.stdout);
流的最佳实践:管道(pipe)
pipe\(\) 是流的核心用法,自动连接可读→转换→可写流,自动处理数据流动和错误:
可读流.pipe(转换流).pipe(可写流);
✅ 优势:无需手动监听 data/end 事件,代码简洁,自动解决大部分背压问题。
二、背压(Back Pressure)问题
1. 什么是背压?
核心现象:读取速度 > 写入速度,数据堆积在内存中,导致内存暴涨、程序崩溃。
举个例子:
-
可读流:每秒读 100MB 数据
-
可写流:每秒只能写 10MB 数据
-
未处理的数据会在内存中排队,这就是背压
2. 背压的危害
-
内存占用飙升,OOM(内存溢出)导致服务宕机
-
系统卡顿,IO 阻塞
三、背压问题的解决方案
Node.js 流内置了背压处理机制,我们只需要正确使用 API 即可:
方案 1:优先使用 pipe\(\) / pipeline\(\)(推荐)
这是最简单、最安全的方案,Node.js 底层自动处理背压:
-
pipe\(\):基础管道,错误会中断流 -
stream\.pipeline\(\):官方推荐,自动清理资源、错误捕获更友好
const { createReadStream, createWriteStream } = require('fs');
const { pipeline } = require('stream');
const { promisify } = require('util');
const pipelineAsync = promisify(pipeline);
// 自动处理背压:读 → 写
async function copyFile() {
try {
await pipelineAsync(
createReadStream('large-file.txt'),
createWriteStream('copy.txt')
);
console.log('拷贝完成,背压自动处理');
} catch (err) {
console.error('出错:', err);
}
}
方案 2:手动处理背压(理解原理)
如果必须手动监听 data 事件,需要判断 writable\.write\(\) 的返回值:
-
返回
false:写入队列已满,停止读取 -
监听
drain事件:写入队列空了,恢复读取
const read = createReadStream('large.txt');
const write = createWriteStream('out.txt');
read.on('data', (chunk) => {
// 写入队列满了,暂停读取
if (!write.write(chunk)) {
read.pause();
}
});
// 写入完成,恢复读取
write.on('drain', () => read.resume());
read.on('end', () => write.end());
方案 3:使用 async iterator(现代写法)
Node.js 支持用 for await\.\.\.of 遍历流,自动处理背压:
async function copy() {
const read = createReadStream('large.txt');
const write = createWriteStream('out.txt');
for await (const chunk of read) {
await write.write(chunk); // 自动等待写入完成,无背压
}
write.end();
}
总结
1. 三大流使用场景
| 流类型 | 定位 | 核心场景 |
|---|---|---|
| Readable | 数据源 | 读文件、网络请求、自定义数据 |
| Writable | 数据目的地 | 写文件、网络响应、数据落盘 |
| Transform | 中间加工层 | 压缩、加密、格式转换、数据清洗 |
2. 背压核心
-
原因:读得快、写得慢,数据堆积内存
-
最优解:用
stream\.pipeline\(\)连接流,底层自动处理背压 -
手动解:
pause\(\)+drain+resume\(\)控制读取速度
最终结论
日常开发 99% 的场景,直接用 pipeline () 就够了,简洁、安全、自动解决背压,这是 Node.js 官方的最佳实践。
(注:文档部分内容可能由 AI 生成)
本文由萧兮的博客原创发布,欢迎转载,转载务必保留原文链接。
萧兮的博客:https://www.20010515.xyz · 原文:https://www.20010515.xyz/posts/019e4ece-aaeb-7161-963f-274a686ff6eb