FORMA

Node.js 流(Stream)详解

流是 Node.js 中处理有序数据的抽象接口,用于以连续方式异步读取或写入数据。它们特别适合处理大文件、网络通信或任何可能无限产生的数据。Node.js 内置的 stream 模块提供了四种基本流类型,以及用于组合和错误处理的工具函数。

一、流的四种基本类型

类型描述方向典型应用
可读流 (Readable)数据可以从中读取单向:流出fs.createReadStreamhttp.IncomingMessage
可写流 (Writable)数据可以写入其中单向:流入fs.createWriteStreamhttp.ServerResponse
双工流 (Duplex)既可读又可写双向独立net.Socketcrypto.createCipheriv
转换流 (Transform)双工流的一种,会对数据进行转换双向且数据经过处理zlib.createGzipcrypto.createHash

示例代码

js
const { Readable, Writable, Duplex, Transform } = require("stream");

// 可读流
const readable = Readable.from(["Hello", "World"]);

// 可写流
const writable = fs.createWriteStream("./output.txt");

// 双工流
const duplex = new Duplex({
    read(size) {
        /* 实现读 */
    },
    write(chunk, encoding, callback) {
        /* 实现写 */
    },
});

// 转换流(大小写转换)
const transform = new Transform({
    transform(chunk, encoding, callback) {
        const output = chunk.toString().toUpperCase();
        callback(null, output);
    },
});

二、背压机制与 pipe / pipeline

1. 背压机制(Backpressure)

当可读流的数据产生速度超过可写流的消费速度时,数据会在内部缓冲区累积。背压机制通过控制数据流动来防止内存溢出:

  • 可读流有一个内部 highWaterMark(默认 16KB)。当缓冲区大小超过该值,readable.push() 返回 false,表示应暂停读取。
  • 可写流也有 highWaterMark,当写入速度慢时,writable.write() 返回 false,并触发 drain 事件。
  • pipepipeline 中,背压自动处理;手动实现时需要检查返回值并监听 drain

2. readable.pipe(writable)

最简单的管道方式,自动处理背压和结束事件。缺点:错误处理不完善(需要单独监听 error 事件)。

js
const readStream = fs.createReadStream("bigfile.txt");
const writeStream = fs.createWriteStream("copy.txt");
readStream.pipe(writeStream);
writeStream.on("finish", () => console.log("复制完成"));

3. stream.pipeline(推荐)

从 Node.js 10 开始引入,支持 Promise 并自动清理流、转发错误,避免内存泄漏。

js
const { pipeline } = require("stream");
const { promisify } = require("util");
const pipelineAsync = promisify(pipeline);

// 或者使用 stream/promises 直接导入 Promise 版本
const { pipeline: pipelinePromise } = require("stream/promises");

async function run() {
    await pipelinePromise(
        fs.createReadStream("input.txt"),
        zlib.createGzip(),
        fs.createWriteStream("input.gz")
    );
    console.log("压缩完成");
}

三、自定义流实现

通过继承核心流类并实现必要的方法,可以创建定制流。

1. 自定义可读流

需要实现 _read(size) 方法,内部调用 push(chunk)push(null) 结束。

js
const { Readable } = require("stream");

class Counter extends Readable {
    constructor(options) {
        super(options);
        this.max = 10;
        this.index = 0;
    }
    _read(size) {
        if (this.index < this.max) {
            this.push(String(this.index++));
        } else {
            this.push(null); // 结束
        }
    }
}

const counter = new Counter();
counter.pipe(process.stdout); // 输出 0 1 2 ... 9

2. 自定义可写流

需要实现 _write(chunk, encoding, callback),处理数据后调用 callback()

js
const { Writable } = require("stream");

class MyWritable extends Writable {
    _write(chunk, encoding, callback) {
        console.log(`写入: ${chunk.toString()}`);
        callback(); // 继续接收下一个块
    }
}

const writable = new MyWritable();
writable.write("Hello\n");
writable.write("World");
writable.end();

3. 自定义转换流

实现 _transform(chunk, encoding, callback),处理后调用 callback(null, transformedChunk)

js
const { Transform } = require("stream");

class UpperCaseTransform extends Transform {
    _transform(chunk, encoding, callback) {
        const upper = chunk.toString().toUpperCase();
        callback(null, upper);
    }
}

const transform = new UpperCaseTransform();
process.stdin.pipe(transform).pipe(process.stdout);

4. 自定义双工流

需同时实现 _read_write,且两者独立。如果不涉及数据转换,直接使用 Duplex

四、实用工具:stream/promises 中的 pipelinefinished

1. pipeline(Promise 版本)

用于串联多个流,返回 Promise,并在任何流发生错误时中止并清理。

js
const { pipeline } = require("stream/promises");
const fs = require("fs");
const zlib = require("zlib");

async function compress() {
    await pipeline(
        fs.createReadStream("file.txt"),
        zlib.createGzip(),
        fs.createWriteStream("file.gz")
    );
    console.log("压缩成功");
}

2. finished

检测流是否完成(或出错),返回 Promise。常用于手动管理流结束。

js
const { finished } = require("stream/promises");

async function monitor(stream) {
    try {
        await finished(stream);
        console.log("流已完成");
    } catch (err) {
        console.error("流错误:", err);
    }
}

const rs = fs.createReadStream("somefile.txt");
monitor(rs);

finished 也可用于每个独立流,而 pipeline 更适合多流串联。

注意stream/promises 从 Node.js 15 开始稳定可用,对于较低版本可使用 util.promisify(require('stream').finished)

五、背压手动处理示例

当不使用 pipepipeline 时,需要手动处理背压:

js
const rs = fs.createReadStream("large.txt");
const ws = fs.createWriteStream("copy.txt");

rs.on("data", chunk => {
    const canWrite = ws.write(chunk);
    if (!canWrite) {
        rs.pause(); // 暂停读取
        ws.once("drain", () => rs.resume()); // 等待可写流排空
    }
});
rs.on("end", () => ws.end());
ws.on("finish", () => console.log("复制结束"));

pipe 内部就是实现这样的逻辑。

六、对象模式(Object Mode)

默认流传输的是 Buffer 或字符串,通过设置 objectMode: true 可以传输任意 JavaScript 对象(如 JSON 数据),但会禁用背压控制,需谨慎使用。

js
const { Readable } = require("stream");
const objStream = new Readable({
    objectMode: true,
    read() {
        this.push({ id: 1 });
        this.push({ id: 2 });
        this.push(null);
    },
});
objStream.on("data", obj => console.log(obj));

总结

流类型核心方法用途
Readable_read产生数据
Writable_write消费数据
Duplex_read + _write双向通信
Transform_transform数据转换(如压缩、加密)

关键最佳实践

  • 使用 pipeline 组合流,自动处理背压和错误。
  • 避免直接使用 pipe 而不监听错误。
  • 自定义流时继承正确的基类,并实现必需方法。
  • 使用 finished 检测流结束状态。

推荐参考资料

Series

core

7 / 7