Node.js Stream 实战:处理大文件与数据管道的正确姿势

By | 2026年6月24日

为什么你需要掌握 Stream?

在 Node.js 中处理大文件(如几 GB 的日志、视频或数据集)时,传统的一次性读取所有内容到内存的方式会轻易耗尽内存,甚至导致进程崩溃。Stream(流)正是为此而生——它将数据分块处理,每次只占用一小块内存,从而高效处理任意大小的数据。

然而,很多开发者对 Stream 的认知停留在“用 pipe 连接”的层面,遇到背压、错误处理、自定义流等问题时常常手足无措。本文将从实战出发,带你深入理解 Stream 的核心机制,并给出可复用的最佳实践。

Stream 基础回顾

Node.js 中有四种基本的流类型:

  • Readable(可读流):数据源,如 fs.createReadStream
  • Writable(可写流):数据目标,如 fs.createWriteStream
  • Transform(转换流):既可读又可写,可在读写过程中修改数据,如 zlib.createGzip
  • Duplex(双工流):既可读又可写,但读写独立,如 net.Socket

核心方法:

  • readable.pipe(writable):将可读流连接到可写流,自动处理背压
  • readable.on('data', chunk => ...):手动消费数据(不推荐,除非需要精细控制)
  • writable.write(chunk):写入数据块
  • writable.end():结束写入

实战一:大文件复制与背压控制

错误示例:一次性读取


const fs = require('fs');

// 危险!大文件会耗尽内存
fs.readFile('bigfile.txt', (err, data) => {
    fs.writeFile('copy.txt', data, err => {
        if (err) console.error(err);
    });
});

正确做法:使用 pipe 自动背压


const fs = require('fs');

const readable = fs.createReadStream('bigfile.txt');
const writable = fs.createWriteStream('copy.txt');

readable.pipe(writable);

readable.on('error', err => console.error('读取错误:', err));
writable.on('error', err => console.error('写入错误:', err));
writable.on('finish', () => console.log('复制完成'));

背压(Backpressure):当写入速度慢于读取速度时,pipe 会自动暂停读取,等待写入完成后再继续。这是 Stream 最强大的特性之一。

手动控制背压(高级场景)

有时你需要更精细的控制,比如统计进度或转换数据。此时可以监听 drain 事件:


const fs = require('fs');

const readable = fs.createReadStream('bigfile.txt');
const writable = fs.createWriteStream('copy.txt');

let totalBytes = 0;

readable.on('data', chunk => {
    totalBytes += chunk.length;
    // 如果 write 返回 false,表示缓冲区已满,应暂停读取
    if (!writable.write(chunk)) {
        readable.pause();
    }
    console.log(`已处理 ${totalBytes} 字节`);
});

writable.on('drain', () => {
    readable.resume();
});

readable.on('end', () => writable.end());
readable.on('error', err => console.error(err));
writable.on('finish', () => console.log('复制完成'));

> 💡 注意:手动实现背压时,务必确保在 drain 事件中恢复读取,否则可能导致内存泄漏或程序挂起。

实战二:自定义 Transform 流实现数据转换

假设我们需要将 CSV 文件中的每一行转成 JSON 对象,并过滤掉某些字段。我们可以创建一个 Transform 流:


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

class CsvToJsonTransform extends Transform {
    constructor(options = {}) {
        super(options);
        this.buffer = ''; // 用于存储不完整的行
    }

    _transform(chunk, encoding, callback) {
        this.buffer += chunk.toString();
        const lines = this.buffer.split('\n');
        // 保留最后一个可能不完整的行
        this.buffer = lines.pop();

        for (const line of lines) {
            if (line.trim() === '') continue;
            const [name, age, city] = line.split(',');
            const json = JSON.stringify({ name, age: parseInt(age), city }) + '\n';
            this.push(json);
        }
        callback();
    }

    _flush(callback) {
        // 处理最后剩余的行
        if (this.buffer.trim() !== '') {
            const [name, age, city] = this.buffer.split(',');
            const json = JSON.stringify({ name, age: parseInt(age), city }) + '\n';
            this.push(json);
        }
        callback();
    }
}

// 使用
const readable = fs.createReadStream('data.csv');
const writable = fs.createWriteStream('data.json');
const transformer = new CsvToJsonTransform();

readable.pipe(transformer).pipe(writable);

readable.on('error', err => console.error(err));
transformer.on('error', err => console.error(err));
writable.on('finish', () => console.log('转换完成'));

> 💡 关键点: > – transform 方法处理每个数据块,注意处理跨块的行(buffer 机制) > – flush 方法在流结束时被调用,用于处理剩余数据 > – 使用 this.push() 输出转换后的数据

实战三:数据管道与错误处理最佳实践

实际项目中,我们经常需要串联多个流(如解压、解密、转换、压缩)。错误处理是容易忽略的环节,一个流的错误如果不处理,会传播到整个管道,导致程序崩溃。

管道中的错误传播


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

const readable = fs.createReadStream('archive.gz');
const gunzip = zlib.createGunzip();
const writable = fs.createWriteStream('output.txt');

// 使用 pipeline 而不是 pipe,自动处理错误和清理
pipeline(
    readable,
    gunzip,
    writable,
    (err) => {
        if (err) {
            console.error('Pipeline 失败:', err);
        } else {
            console.log('Pipeline 成功');
        }
    }
);

为什么推荐 pipeline

  • pipe 不会自动销毁流,错误可能导致内存泄漏
  • pipeline 在 Node.js 10+ 中引入,会自动销毁所有流并调用回调
  • 手动使用 pipe 需要监听每个流的错误并手动销毁,非常繁琐

自定义错误处理与重试

对于网络请求等场景,我们可以封装一个可重试的流:


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

class RetryTransform extends Transform {
    constructor(options = {}) {
        super(options);
        this.retries = options.retries || 3;
    }

    _transform(chunk, encoding, callback) {
        // 模拟一个可能失败的操作
        this.processWithRetry(chunk, 0, callback);
    }

    processWithRetry(chunk, attempt, callback) {
        try {
            // 假设这里可能抛出错误
            if (Math.random() < 0.3) throw new Error('随机失败');
            const result = chunk.toString().toUpperCase();
            this.push(result);
            callback();
        } catch (err) {
            if (attempt < this.retries) {
                console.log(`重试第 ${attempt + 1} 次`);
                setTimeout(() => this.processWithRetry(chunk, attempt + 1, callback), 1000);
            } else {
                callback(err); // 超过重试次数,传递错误
            }
        }
    }
}

> ⚠️ 注意:重试逻辑应谨慎使用,避免无限重试导致资源耗尽。

实战四:对象模式流(Object Mode)

默认情况下,流处理的是 Buffer 或字符串。但有时我们需要处理 JavaScript 对象,这时需要启用对象模式:


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

const objectStream = new Transform({
    objectMode: true,
    transform(chunk, encoding, callback) {
        // chunk 是一个对象
        chunk.timestamp = Date.now();
        this.push(chunk);
        callback();
    }
});

objectStream.write({ name: 'Alice' });
objectStream.on('data', obj => console.log(obj));
// 输出: { name: 'Alice', timestamp: 1700000000000 }

> 💡 注意:对象模式会大幅降低性能,因为每个对象都是一个独立的数据块,失去了 Buffer 的批处理优势。仅当需要处理结构化数据且数据量不大时使用。

性能优化与监控

高水位线(HighWaterMark)

highWaterMark 控制流内部缓冲区的最大字节数(或对象数)。调整它可以影响背压行为:


// 增大缓冲区以减少背压触发频率(适合高速读取)
const readable = fs.createReadStream('bigfile.txt', { highWaterMark: 1024 * 1024 }); // 1MB
// 减小缓冲区以降低内存占用(适合低速写入)
const writable = fs.createWriteStream('output.txt', { highWaterMark: 16 * 1024 }); // 16KB

使用 perf_hooks 监控流性能


const { performance, PerformanceObserver } = require('perf_hooks');
const fs = require('fs');

const obs = new PerformanceObserver((items) => {
    console.log(items.getEntries()[0]);
});
obs.observe({ entryTypes: ['measure'] });

performance.mark('start');

const readable = fs.createReadStream('bigfile.txt');
const writable = fs.createWriteStream('copy.txt');

readable.pipe(writable);

writable.on('finish', () => {
    performance.mark('end');
    performance.measure('copy', 'start', 'end');
});

总结

  • 永远使用 Stream 处理大文件,避免 readFile/writeFile
  • 优先使用 pipeline 而不是 pipe,自动处理错误和清理
  • 自定义 Transform 流时,注意跨块数据处理和 _flush
  • 理解背压机制,手动控制时务必监听 drain 事件
  • 对象模式只在小数据量时使用,注意性能开销
  • 调整 highWaterMark 可以优化特定场景的性能

延伸阅读

  • Node.js 官方文档:Stream
  • 使用 stream.finished 检测流是否结束
  • 考虑使用 pump 库(pipeline 的前身)兼容旧版本
  • 探索 readable.pipe(destination, { end: false }) 实现多路复用