Node.js: Stream 流
最后更新:2026-08-26
Charlie 接到一个紧急任务:分析服务器上 2GB 的访问日志。他用 fs.readFile 一加载,内存直接爆了——OOM 崩溃。组长说:"文件再大也不能一次吞,用 Stream 一口一口吃!"Charlie 改用 createReadStream 逐行处理,内存始终稳定在 50MB。从此他逢人就说:Stream 就是 Node.js 处理大数据的瑞士军刀。
1. Stream 基础概念
Stream 是 Node.js 中处理流式数据的抽象接口,数据像水流一样逐块(chunk)传输,而非一次性加载到内存。
- 核心思想:数据分块处理,无需等待全部数据就位即可开始操作
- 继承关系:所有流都是
EventEmitter的实例,通过事件通知状态变化 - 适用场景:大文件读写、网络传输、数据压缩与转换
- 内置模块:
fs、zlib、crypto、http等均提供了流接口 - 全局对象:
stream模块提供Readable、Writable、Duplex、Transform基类
(1) 四种流类型
| 类型 | 说明 | 典型场景 | 输入 | 输出 |
|---|---|---|---|---|
| Readable | 可读流,数据源头 | fs.createReadStream、process.stdin |
无 | 有 |
| Writable | 可写流,数据终点 | fs.createWriteStream、process.stdout |
有 | 无 |
| Duplex | 双工流,可读可写(独立) | net.Socket、tls.Socket |
有 | 有 |
| Transform | 转换流,输出是输入的变换 | zlib.createGzip、crypto.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();
data事件在流动模式下自动触发end与finish的区别:end表示读完成,finish表示写完成error事件必须监听,否则未捕获的异常会导致进程崩溃close事件不是所有流都会触发,取决于底层实现
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 数据流图
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
pipe()自动管理数据流节奏,可读端在可写端繁忙时暂停pipe()不处理错误,需对每个流单独监听error事件- 使用
stream.pipeline()可自动清理和错误传播,比pipe()更安全
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} 字节`);
});
- 大文件必须用流,避免 OOM
highWaterMark越小内存越省,但系统调用越频繁start/end选项可实现文件切片读取
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();
});
write()返回false表示内部缓冲区已满,应暂停读取drain事件表示缓冲区已排空,可以恢复写入- 忽略背压会导致内存持续增长,最终 OOM
pipeline()比pipe()更推荐,自动处理错误传播与资源清理
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'));
_transform(chunk, encoding, callback)中,callback(null, data)推送变换结果callback(err)可传递错误zlib.createGzip()/createGunzip()是最常用的内置 Transform 流- 自定义 Transform 流时也可实现
_flush(callback)处理剩余数据
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('文件读取完毕');
});
readline.createInterface将 Readable 流包装为逐行读取接口crlfDelay: Infinity可正确处理\r\n换行pause()后不再触发data事件,直到调用resume()pipe()会将流切换到流动模式
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);
});
- 管道中每个环节都是流,数据逐块流过,内存占用极低
pipe()自动在相邻流间处理背压- 错误需在每个流上单独监听,或使用
pipeline()统一处理 - Transform 流
filterError只保留包含ERROR的行
❓ 常见问题
Q 什么是流(Stream)?
A 流是一种处理数据的抽象接口,数据可以逐块读取或写入,而不需要一次性加载到内存中。
Q Readable 和 Writable 流有什么区别?
A Readable 流是数据源,可以从中读取数据;Writable 流是数据目标,可以向其中写入数据。
Q 什么时候应该使用流?
A 处理大文件、网络传输、实时数据处理等场景,数据量较大或无法一次性加载到内存时。
Q pipe 方法的作用是什么?
A pipe 将可读流的输出连接到可写流的输入,自动管理数据流动和背压处理。
Q 什么是背压(Backpressure)?
A 当可写流处理数据的速度慢于可读流产生数据的速度时,可写流会向可读流发送信号暂停读取,这就是背压机制。
- 什么时候用 Stream? 当数据量超过可用内存、需要边读边处理、或处理实时数据源时使用 Stream;小文件用 readFile 更简单。
- pipe() 自动处理背压吗? 是的,
pipe()内部检测可写端的write()返回值,自动调用pause()/resume()管理流速。 - 如何逐行读取文件? 使用
readline.createInterface({ input: createReadStream(path) }),监听line事件即可逐行获取。 - Transform 和 Duplex 有什么区别? Duplex 的读写独立无关,Transform 的输出由输入经变换产生,只需实现
_transform()方法。 - 流的错误如何处理? 对每个流监听
error事件,或使用stream.pipeline()自动传播错误并清理资源,避免未捕获异常导致进程崩溃。 - highWaterMark 设多大合适? 默认 64KB 适合多数场景;处理大块二进制可增大到 256KB,内存受限时减小到 16KB。
- pipe() 和 pipeline() 选哪个? 推荐用
stream.pipeline(),它能自动传播错误、清理资源;pipe()不处理错误且不自动关闭流。
📖 小节
- Stream 基础概念的核心概念与使用方法
- 流事件详解的核心概念与使用方法
- pipe() 链式连接的核心概念与使用方法
- fs 流与文件操作的核心概念与使用方法
- 背压(Backpressure)机制的核心概念与使用方法
- Transform 流的核心概念与使用方法
- 流的暂停与恢复的核心概念与使用方法
- 综合示例:日志处理管道的核心概念与使用方法
📝 作业
- 完成本课所有代码示例,确保每个示例都能正确运行
- 修改综合示例,添加自己的扩展功能
- 查阅官方文档,找出本课未涉及的1-2个API并编写测试代码
- 思考:在实际项目中,你会如何应用本课学到的知识?
- 尝试将本课知识与前面课程的内容结合,构建一个小项目