Stream 流
本教程共 76 篇 · 第 21 篇 · 更新于 2026-07-25 · 约 6 分钟阅读
21. Stream 流
本节目标:四种流类型、pipe 管道和背压(backpressure)机制。
想象你要把一盆水从 A 点运到 B 点。第一种办法是等盆装满,再整个端过去;第二种是接一根水管,让水一边流一边处理。Stream 就是 Node.js 里的那根水管。
读一个大文件、接收 HTTP 请求体、压缩日志、处理视频流——这些场景下,Stream 能把内存占用压得很低,而且处理延迟更小,数据还没到齐就可以开始干活。
21.1 为什么需要 Stream
先做个对比。下面这段代码读一个 1GB 的视频文件:
import { readFile } from 'node:fs/promises';
// 一次性把整个文件塞进内存
const data = await readFile('movie.mp4');
console.log(data.length);
如果服务器只有 512MB 可用内存,这行代码直接让进程崩溃。换成 Stream:
import { createReadStream } from 'node:fs';
const stream = createReadStream('movie.mp4', { highWaterMark: 64 * 1024 });
stream.on('data', chunk => {
console.log('收到一块:', chunk.length);
// 64KB 一块一块处理
});
内存里始终只驻留 64KB,1GB 的文件也能平稳处理。
21.2 四种流类型
Node.js 把流分成四类:
| 类型 | 角色 | 例子 |
|---|---|---|
| Readable | 数据源(读) | fs.createReadStream、http.IncomingMessage |
| Writable | 数据目的地(写) | fs.createWriteStream、process.stdout |
| Duplex | 可读可写 | TCP Socket |
| Transform | 读写中间加工数据 | zlib.createGzip、crypto.createCipher |
所有流都继承自 EventEmitter,靠事件驱动数据流动。
21.3 Readable:读数据
流动模式(flowing)
监听 data 事件,数据会自动推给你:
import { createReadStream } from 'node:fs';
const reader = createReadStream('log.txt', { encoding: 'utf8' });
reader.on('data', chunk => {
console.log('收到:', chunk);
});
reader.on('end', () => {
console.log('读完了');
});
reader.on('error', err => {
console.error('出错:', err);
});
暂停模式(paused)
用 readable 事件配合 read(),主动权在你手里:
import { createReadStream } from 'node:fs';
const reader = createReadStream('log.txt', { encoding: 'utf8' });
reader.on('readable', () => {
let chunk;
while ((chunk = reader.read()) !== null) {
console.log('手动读取:', chunk);
}
});
日常写业务代码,data 事件更省事;需要精细控制消费速率时,再用暂停模式。
异步迭代器(推荐)
v10 以后,Readable 流实现了异步迭代协议,写法最干净:
import { createReadStream } from 'node:fs';
const reader = createReadStream('log.txt', { encoding: 'utf8' });
for await (const chunk of reader) {
console.log('迭代:', chunk);
}
for await...of 自动处理背压和错误传播,比事件监听少写很多样板代码。
21.4 Writable:写数据
import { createWriteStream } from 'node:fs';
const writer = createWriteStream('output.txt', { encoding: 'utf8' });
writer.write('第一行\n');
writer.write('第二行\n');
writer.end('结束\n');
writer.on('finish', () => {
console.log('全部写入完成');
});
write() 返回一个布尔值,表示内部缓冲区是否已满。这个返回值是背压机制的关键。
21.5 pipe:连接流
readable.pipe(writable) 是最经典的流用法,像接水管一样把可读流和可写流连起来:
import { createReadStream, createWriteStream } from 'node:fs';
import { createGzip } from 'node:zlib';
// 读取 -> 压缩 -> 写入
const source = createReadStream('big.log');
const gzip = createGzip();
const dest = createWriteStream('big.log.gz');
source.pipe(gzip).pipe(dest);
dest.on('finish', () => {
console.log('压缩完成');
});
.pipe() 会自动处理背压:如果写入端忙不过来,读取端会暂停;等写入端清空缓冲区,再恢复读取。
但 .pipe() 有个老毛病——错误不会自动传播。如果 gzip 中间报错,source 和 dest 不会自动关闭,可能导致内存泄漏或文件句柄没释放。
21.6 pipeline:更安全的管道
import { pipeline } from 'node:stream/promises';
import { createReadStream, createWriteStream } from 'node:fs';
import { createGzip } from 'node:zlib';
await pipeline(
createReadStream('big.log'),
createGzip(),
createWriteStream('big.log.gz')
);
console.log('pipeline 完成');
stream/promises 里的 pipeline 返回 Promise,内部会自动:
- 把错误从中间流抛到外部。
- 出错时关闭所有参与的流。
- 最后一个流正常结束时 resolve。
新代码优先用 pipeline,别再用 .pipe() 了。
21.7 Backpressure:背压
背压是 Stream 的核心机制,也是很多人用流时踩坑的地方。
问题场景
读取端很快,写入端很慢。如果不做限制,读取端拼命往内存里塞数据,写入端处理不过来,内存会被撑爆。
解决原理
Writable 的 write(chunk) 在内部缓冲区达到 highWaterMark(默认 16KB)时,会返回 false:
const writer = createWriteStream('slow.txt');
const ok = writer.write(someData);
if (!ok) {
console.log('缓冲区满了,别写了');
}
Readable 收到信号后会自动暂停 data 事件。等 Writable 把缓冲区消化完,发出 drain 事件,Readable 再恢复:
function writeLots(writer, data) {
let i = 0;
function write() {
let ok = true;
do {
ok = writer.write(data[i++]);
} while (i < data.length && ok);
if (i < data.length) {
// 等 drain 事件再续写
writer.once('drain', write);
}
}
write();
}
手动控制背压
如果你自己实现自定义 Writable,必须在 _write 里调用 callback 告诉底层「我处理完了」:
import { Writable } from 'node:stream';
const slowWriter = new Writable({
write(chunk, encoding, callback) {
// 模拟慢速写入,比如发到远程服务器
setTimeout(() => {
console.log('写入:', chunk.toString());
callback(); // 通知可以接收下一块了
}, 100);
}
});
不调用 callback,Stream 就会认为你还在处理,不会继续推送数据。
21.8 Transform:中间加工
Transform 流是 Duplex 的一种,读进去的数据和写出去的数据不一样,中间可以做转换。
import { Transform } from 'node:stream';
import { createReadStream, createWriteStream } from 'node:fs';
import { pipeline } from 'node:stream/promises';
const upperCase = new Transform({
transform(chunk, encoding, callback) {
this.push(chunk.toString().toUpperCase());
callback();
}
});
await pipeline(
createReadStream('input.txt'),
upperCase,
createWriteStream('output.txt')
);
常见 Transform 实现:JSON 行解析、CSV 字段过滤、加密解密、压缩解压。
21.9 Object Mode:流传对象
默认情况下流里流的是 Buffer 或 string。设置 objectMode: true 后,可以流传任意 JavaScript 对象:
import { Readable, Transform } from 'node:stream';
import { pipeline } from 'node:stream/promises';
const objectSource = Readable.from([
{ id: 1, name: 'Alice' },
{ id: 2, name: 'Bob' },
{ id: 3, name: 'Charlie' }
]);
const filter = new Transform({
objectMode: true,
transform(obj, encoding, callback) {
if (obj.id > 1) this.push(obj);
callback();
}
});
for await (const obj of objectSource.pipe(filter)) {
console.log(obj);
}
Readable.from(array) 是个快捷方法,把数组包装成可读流。
21.10 Web Stream 互操作
Node.js 内置了 WHATWG Streams 标准(ReadableStream、WritableStream、TransformStream),和原生的 fetch 返回的 body 是同一套接口:
const response = await fetch('https://example.com/large.bin');
// fetch 返回的 body 是 Web ReadableStream
// 可以用 pipeline 直接接到 Node.js 的写入流
import { createWriteStream } from 'node:fs';
import { pipeline } from 'node:stream/promises';
import { Readable } from 'node:stream';
await pipeline(
Readable.fromWeb(response.body),
createWriteStream('downloaded.bin')
);
Readable.fromWeb 和 Readable.toWeb 用来在两种流标准之间转换。