首页 / Node.js 教程 / Stream 流

Node.js 教程

Stream 流

本教程共 76 篇 · 第 21 篇 · 更新于 2026-07-25 · 约 6 分钟阅读

Node.jsStreampipe背压

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.createReadStreamhttp.IncomingMessage
Writable数据目的地(写)fs.createWriteStreamprocess.stdout
Duplex可读可写TCP Socket
Transform读写中间加工数据zlib.createGzipcrypto.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 中间报错,sourcedest 不会自动关闭,可能导致内存泄漏或文件句柄没释放。

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,内部会自动:

  1. 把错误从中间流抛到外部。
  2. 出错时关闭所有参与的流。
  3. 最后一个流正常结束时 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:流传对象

默认情况下流里流的是 Bufferstring。设置 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 标准(ReadableStreamWritableStreamTransformStream),和原生的 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.fromWebReadable.toWeb 用来在两种流标准之间转换。