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 可读流

定义数据源,负责产生 / 读取数据,向外流出数据。

  • 核心:你可以从它里面读取数据

使用场景

  1. 读取大文件(日志、视频、CSV),不占用大量内存

  2. 网络请求接收数据(HTTP 请求体、WebSocket 消息)

  3. 读取数据库 / 接口返回的流式数据

  4. 自定义数据源(生成随机数、实时日志推送)

简单示例

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 可写流

定义数据目的地,负责接收 / 写入数据,向内流入数据。

  • 核心:你可以向它写入数据

使用场景

  1. 写入大文件(日志落盘、导出大数据)

  2. 网络响应发送数据(HTTP 响应、文件下载)

  3. 写入数据库 / 消息队列(流式写入)

简单示例

const { createWriteStream } = require('fs');
const writeStream = createWriteStream('output.txt');

writeStream.write('第一行数据\n');
writeStream.write('第二行数据\n');
writeStream.end('最后一行数据'); // 结束写入

3. Transform 转换流

定义中间处理层可读 + 可写,数据流经它时被修改 / 处理,再输出。

  • 核心:读 → 处理 → 写,是数据的 “加工管道”。

使用场景

  1. 数据压缩 / 解压(gzip)

  2. 数据加密 / 解密

  3. 文本替换、格式转换(CSV 转 JSON、大写转小写)

  4. 日志过滤、数据清洗

简单示例

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