Node.js: Stream 流

最后更新:2026-08-26

Charlie 接到一个紧急任务:分析服务器上 2GB 的访问日志。他用 fs.readFile 一加载,内存直接爆了——OOM 崩溃。组长说:"文件再大也不能一次吞,用 Stream 一口一口吃!"Charlie 改用 createReadStream 逐行处理,内存始终稳定在 50MB。从此他逢人就说:Stream 就是 Node.js 处理大数据的瑞士军刀。

1. Stream 基础概念

Stream 是 Node.js 中处理流式数据的抽象接口,数据像水流一样逐块(chunk)传输,而非一次性加载到内存。

(1) 四种流类型

类型 说明 典型场景 输入 输出
Readable 可读流,数据源头 fs.createReadStreamprocess.stdin
Writable 可写流,数据终点 fs.createWriteStreamprocess.stdout
Duplex 双工流,可读可写(独立) net.Sockettls.Socket
Transform 转换流,输出是输入的变换 zlib.createGzipcrypto.createCipheriv

(2) 流的两种模式

Readable 流有两种运行模式:

特性 流动模式(Flowing) 暂停模式(Paused)
数据获取 自动推送,通过 data 事件消费 需手动调用 read() 拉取
触发方式 添加 data 监听 / 调用 pipe() / 调用 resume() 初始默认状态 / 调用 pause()
背压处理 pipe() 自动处理;手动需监听 drain 由消费者控制节奏,天然背压
适用场景 高吞吐、持续数据源 需要精确控制读取时机


2. 流事件详解

所有流都基于 EventEmitter,通过事件机制传递数据与状态。

(1) 通用事件

事件 触发时机 适用流 回调参数
data 收到新 chunk Readable chunk
end 数据读取完毕 Readable
error 发生错误 所有流 Error
close 流底层资源关闭 所有流
finish end() 后所有数据刷完 Writable

▶ 示例:监听 Readable 流事件

JAVASCRIPT
const fs = require('fs');
const rs = fs.createReadStream('./access.log', { highWaterMark: 64 * 1024 });

rs.on('data', (chunk) => {
  console.log(`收到 ${chunk.length} 字节`);
});

rs.on('end', () => {
  console.log('读取完毕');
});

rs.on('error', (err) => {
  console.error('流出错:', err.message);
});
▶ 试一试

▶ 示例:监听 Writable 流事件

JAVASCRIPT
const ws = fs.createWriteStream('./output.txt');

ws.on('finish', () => {
  console.log('所有数据已写入');
});

ws.on('error', (err) => {
  console.error('写入错误:', err.message);
});

ws.write('Hello Stream');
ws.end();
▶ 试一试

3. pipe() 链式连接

pipe() 是 Stream 最强大的方法,它将可读流的输出自动连接到可写流的输入,并自动处理背压。

▶ 示例:(1) pipe 基本用法

JAVASCRIPT
readable.pipe(writable);
▶ 试一试

pipe() 返回目标流,因此可以链式调用:

JAVASCRIPT
readable.pipe(transform1).pipe(transform2).pipe(writable);

▶ 示例:文件复制

JAVASCRIPT
const fs = require('fs');
fs.createReadStream('./source.txt')
  .pipe(fs.createWriteStream('./dest.txt'));
▶ 试一试

▶ 示例:HTTP 响应流式传输

JAVASCRIPT
const http = require('http');
const fs = require('fs');

http.createServer((req, res) => {
  res.writeHead(200, { 'Content-Type': 'text/plain' });
  fs.createReadStream('./large.txt').pipe(res);
}).listen(3000);
▶ 试一试

▶ 示例:(2) Stream pipe 数据流图

100%
flowchart LR
    A[Readable<br/>数据源头] -->|chunk| B[Transform<br/>数据变换]
    B -->|chunk| C[Writable<br/>数据终点]
    C -.->|backpressure| A
    B -.->|backpressure| A

    style A fill:#4CAF50,color:#fff
    style B fill:#FF9800,color:#fff
    style C fill:#2196F3,color:#fff


4. fs 流与文件操作

fs 模块提供了文件系统相关的可读流与可写流。

(1) createReadStream / createWriteStream

选项 说明 默认值
highWaterMark 缓冲区大小(字节) Readable: 64KB / Writable: 16KB
encoding 编码方式 null(Buffer)
start 起始字节位置 0
end 结束字节位置(含) Infinity
flags 文件打开标志 Readable: r / Writable: w

▶ 示例:指定范围读取文件片段

JAVASCRIPT
const fs = require('fs');
const rs = fs.createReadStream('./big.bin', {
  start: 100,
  end: 199,
  highWaterMark: 32
});

rs.on('data', (chunk) => {
  console.log(chunk.length);
});
▶ 试一试

(2) readFile vs createReadStream 对比

对比项 readFile createReadStream
内存占用 文件全部加载到内存 仅占用 highWaterMark 大小
启动延迟 等文件全部读完才回调 立即返回流,边读边处理
适用文件大小 小文件(<10MB) 大文件或持续数据
错误处理 回调中获取 Error 监听 error 事件
可暂停/恢复 不支持 pause() / resume()

▶ 示例:大文件逐块处理

JAVASCRIPT
const fs = require('fs');
let totalBytes = 0;

const rs = fs.createReadStream('./2gb.log');
rs.on('data', (chunk) => {
  totalBytes += chunk.length;
});
rs.on('end', () => {
  console.log(`总计 ${totalBytes} 字节`);
});
▶ 试一试

5. 背压(Backpressure)机制

背压是流控的核心:当可写端处理速度跟不上可读端推送速度时,向可写端施加"反压",让可读端暂停,避免内存堆积。

(1) pipe() 自动处理背压

使用 pipe() 时,背压由 Node.js 内部自动管理,无需手动干预。

(2) 手动处理背压

不使用 pipe() 时,需手动判断 write() 返回值:

JAVASCRIPT
const fs = require('fs');
const rs = fs.createReadStream('./source.txt');
const ws = fs.createWriteStream('./dest.txt');

rs.on('data', (chunk) => {
  const canContinue = ws.write(chunk);
  if (!canContinue) {
    rs.pause();
    ws.once('drain', () => {
      rs.resume();
    });
  }
});

rs.on('end', () => {
  ws.end();
});

▶ 示例:观察背压效果

JAVASCRIPT
const rs = fs.createReadStream('./big.log', { highWaterMark: 1024 });
const ws = fs.createWriteStream('./out.log', { highWaterMark: 512 });

let paused = 0;
rs.on('data', (chunk) => {
  const ok = ws.write(chunk);
  if (!ok) {
    paused++;
    rs.pause();
    ws.once('drain', () => rs.resume());
  }
});
rs.on('end', () => {
  console.log(`背压触发 ${paused} 次`);
  ws.end();
});
▶ 试一试

6. Transform 流

Transform 流是 Duplex 的子类,输出是输入的变换结果,常用于数据压缩、加密、格式转换。

(1) Transform 与 Duplex 的区别

对比项 Duplex Transform
输入输出关系 独立,互不影响 输出由输入经变换产生
必须实现方法 _read() + _write() _transform()
典型用途 网络套接字 压缩、加密、数据转换
内部缓冲 读写各一份 变换中间态

▶ 示例:自定义 Transform 流(大小写转换)

JAVASCRIPT
const { Transform } = require('stream');

const upper = new Transform({
  transform(chunk, encoding, callback) {
    callback(null, chunk.toString().toUpperCase());
  }
});

process.stdin.pipe(upper).pipe(process.stdout);
▶ 试一试

▶ 示例:zlib 压缩文件

JAVASCRIPT
const fs = require('fs');
const zlib = require('zlib');

fs.createReadStream('./access.log')
  .pipe(zlib.createGzip())
  .pipe(fs.createWriteStream('./access.log.gz'));
▶ 试一试

▶ 示例:zlib 解压文件

JAVASCRIPT
const fs = require('fs');
const zlib = require('zlib');

fs.createReadStream('./access.log.gz')
  .pipe(zlib.createGunzip())
  .pipe(fs.createWriteStream('./access_restored.log'));
▶ 试一试

7. 流的暂停与恢复

Readable 流默认处于暂停模式,可通过 pause()resume() 切换。

▶ 示例:暂停与恢复控制

JAVASCRIPT
const fs = require('fs');
const rs = fs.createReadStream('./big.log');
let count = 0;

rs.on('data', (chunk) => {
  count++;
  if (count % 10 === 0) {
    rs.pause();
    console.log(`已处理 ${count} 块,暂停 1 秒`);
    setTimeout(() => rs.resume(), 1000);
  }
});
▶ 试一试

▶ 示例:逐行读取文件(readline)

JAVASCRIPT
const fs = require('fs');
const readline = require('readline');

const rl = readline.createInterface({
  input: fs.createReadStream('./access.log'),
  crlfDelay: Infinity
});

rl.on('line', (line) => {
  if (line.includes('ERROR')) {
    console.log(line);
  }
});

rl.on('close', () => {
  console.log('文件读取完毕');
});
▶ 试一试

8. 综合示例:日志处理管道

构建完整的数据处理管道:读取大日志文件 → 逐行解析 → 过滤错误行 → 压缩输出。

JAVASCRIPT
const fs = require('fs');
const zlib = require('zlib');
const { Transform } = require('stream');

const filterError = new Transform({
  transform(chunk, encoding, callback) {
    const lines = chunk.toString().split('\n');
    const errors = lines.filter(l => l.includes('ERROR')).join('\n');
    callback(null, errors ? errors + '\n' : '');
  }
});

const src = fs.createReadStream('./app.log');
const dest = fs.createWriteStream('./errors.log.gz');

src
  .pipe(filterError)
  .pipe(zlib.createGzip())
  .pipe(dest);

dest.on('finish', () => {
  console.log('错误日志已压缩写入');
});

src.on('error', (err) => {
  console.error('读取失败:', err.message);
});

❓ 常见问题

Q 什么是流(Stream)?
A 流是一种处理数据的抽象接口,数据可以逐块读取或写入,而不需要一次性加载到内存中。
Q Readable 和 Writable 流有什么区别?
A Readable 流是数据源,可以从中读取数据;Writable 流是数据目标,可以向其中写入数据。
Q 什么时候应该使用流?
A 处理大文件、网络传输、实时数据处理等场景,数据量较大或无法一次性加载到内存时。
Q pipe 方法的作用是什么?
A pipe 将可读流的输出连接到可写流的输入,自动管理数据流动和背压处理。
Q 什么是背压(Backpressure)?
A 当可写流处理数据的速度慢于可读流产生数据的速度时,可写流会向可读流发送信号暂停读取,这就是背压机制。

📖 小节


📝 作业

  1. 完成本课所有代码示例,确保每个示例都能正确运行
  2. 修改综合示例,添加自己的扩展功能
  3. 查阅官方文档,找出本课未涉及的1-2个API并编写测试代码
  4. 思考:在实际项目中,你会如何应用本课学到的知识?
  5. 尝试将本课知识与前面课程的内容结合,构建一个小项目
Web-Tutorial.com

Web-Tutorial 技术团队

由多位开发者共同维护的编程教程平台。每篇教程由对应领域的开发者编写和审核,确保内容准确可靠。如发现任何问题,欢迎向我们反馈。

100%

🙏 帮我们做得更好

我们是刚上线的编程教程站,几个人的小团队,精力有限。页面虽经检查,难免还有疏漏——链接失效、排版错乱、内容有误、语言生硬……

如果您发现了,麻烦告诉我们,我们会在收到反馈后第一时间进行修复,再次感谢您的光临 🙏