Node.js Stream 实战:处理大文件与数据管道

By | 2026年7月28日

引言

在 Node.js 开发中,处理大文件或大量数据时,传统的一次性读取方式(如 fs.readFile)往往会导致内存溢出或性能瓶颈。流(Stream)是 Node.js 中处理数据的核心抽象,它允许你以分块的方式处理数据,从而节省内存并提高效率。本文将带你从零开始,掌握 Node.js Stream 的实战技巧,包括读写流、转换流、管道操作以及背压处理。

1. Stream 基础

Node.js 的流有四种类型:

  • Readable:可读流(如 fs.createReadStream
  • Writable:可写流(如 fs.createWriteStream
  • Transform:转换流(如 zlib.createGzip
  • Duplex:双工流(可读可写,如 net.Socket

所有流都是 EventEmitter 的实例,通过事件驱动。

1.1 创建简单的读写流


const fs = require('fs');

// 创建可读流
const readableStream = fs.createReadStream('input.txt', {
  highWaterMark: 16 * 1024, // 每次读取 16KB,默认 64KB
  encoding: 'utf8'
});

// 创建可写流
const writableStream = fs.createWriteStream('output.txt');

// 监听数据事件
readableStream.on('data', (chunk) => {
  console.log(`Received ${chunk.length} bytes of data.`);
  writableStream.write(chunk);
});

readableStream.on('end', () => {
  console.log('Reading finished.');
  writableStream.end();
});

// 错误处理
readableStream.on('error', (err) => console.error('Read error:', err));
writableStream.on('error', (err) => console.error('Write error:', err));

注意:手动管理读写容易出错,比如忘记调用 writableStream.end() 会导致文件不完整。推荐使用 pipe 方法。

2. 管道(Pipe)操作

pipe 是流最强大的功能,它自动管理背压(backpressure),防止数据积压导致内存溢出。


const fs = require('fs');

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

readable.pipe(writable);

writable.on('finish', () => console.log('File copied successfully.'));

2.1 链式管道

管道可以链式调用,形成数据处理流水线。


const fs = require('fs');
const zlib = require('zlib');

// 压缩文件
fs.createReadStream('input.txt')
  .pipe(zlib.createGzip())
  .pipe(fs.createWriteStream('input.txt.gz'))
  .on('finish', () => console.log('Compression done.'));

💡 最佳实践:始终在最后一个流上监听 finishclose 事件,确保所有数据已处理完毕。

3. 自定义转换流

转换流用于在数据传递过程中修改数据。例如,创建一个将文本转为大写的流。


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

class UpperCaseTransform extends Transform {
  _transform(chunk, encoding, callback) {
    // chunk 是 Buffer 或字符串,取决于 encoding
    this.push(chunk.toString().toUpperCase());
    callback();
  }
}

// 使用
const upperCase = new UpperCaseTransform();

process.stdin.pipe(upperCase).pipe(process.stdout);

3.1 实战:CSV 数据清洗

假设我们需要读取一个大型 CSV 文件,过滤掉某些行,并输出为 JSON。


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

// 创建一个过滤转换流
class FilterTransform extends Transform {
  constructor(options) {
    super({ objectMode: true, ...options }); // 使用对象模式
    this.header = null;
  }

  _transform(chunk, encoding, callback) {
    // 假设 chunk 是一行字符串
    const line = chunk.toString().trim();
    if (!line) return callback();

    if (!this.header) {
      this.header = line.split(',');
      return callback();
    }

    const values = line.split(',');
    if (values[1] === 'valid') { // 假设第二列是状态
      const obj = {};
      this.header.forEach((key, i) => { obj[key] = values[i]; });
      this.push(JSON.stringify(obj) + '\n');
    }
    callback();
  }
}

// 使用
const readStream = fs.createReadStream('data.csv');
const writeStream = fs.createWriteStream('output.json');

readStream
  .pipe(new FilterTransform())
  .pipe(writeStream)
  .on('finish', () => console.log('CSV processed.'));

注意:使用 objectMode 后,流处理的是对象而非 Buffer,但要注意内存管理。

4. 背压(Backpressure)详解

背压是指当写入速度慢于读取速度时,流自动控制数据流动的机制。pipe 方法自动处理背压,但手动实现时需注意。

4.1 手动处理背压


const fs = require('fs');

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

readable.on('data', (chunk) => {
  const canContinue = writable.write(chunk);
  if (!canContinue) {
    console.log('Backpressure detected, pausing read.');
    readable.pause();
    writable.once('drain', () => {
      console.log('Drain event, resuming read.');
      readable.resume();
    });
  }
});

readable.on('end', () => writable.end());

💡 最佳实践:除非有特殊需求,否则始终使用 pipepipeline(Node.js 10+)来处理背压。

5. 使用 pipeline 替代 pipe

pipeline 是 Node.js 10 引入的模块化方法,它自动处理错误和清理。


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

pipeline(
  fs.createReadStream('input.txt'),
  zlib.createGzip(),
  fs.createWriteStream('input.txt.gz'),
  (err) => {
    if (err) {
      console.error('Pipeline failed:', err);
    } else {
      console.log('Pipeline succeeded.');
    }
  }
);

pipeline 在发生错误时会自动销毁所有流,避免资源泄漏。

6. 实战:大文件行处理

处理大文件时,逐行读取比按 chunk 读取更方便。但 Node.js 没有内置的行读取流,我们可以自己实现或使用 readline 模块。


const fs = require('fs');
const readline = require('readline');

async function processLargeFile(filePath) {
  const fileStream = fs.createReadStream(filePath);

  const rl = readline.createInterface({
    input: fileStream,
    crlfDelay: Infinity
  });

  let lineCount = 0;
  for await (const line of rl) {
    // 处理每一行
    lineCount++;
    if (lineCount % 100000 === 0) {
      console.log(`Processed ${lineCount} lines`);
    }
  }
  console.log(`Total lines: ${lineCount}`);
}

processLargeFile('huge.log').catch(console.error);

注意readline 内部使用了流,但它是按行分割的,适合处理文本文件。

7. 常见陷阱与最佳实践

7.1 不要忘记错误处理

流事件中的错误必须被捕获,否则会导致进程崩溃。


const fs = require('fs');
const stream = fs.createReadStream('nonexistent.txt');
stream.on('error', (err) => console.error('Error:', err.message));

7.2 避免在 ‘data’ 事件中执行同步操作

同步操作会阻塞事件循环,影响性能。

7.3 使用 highWaterMark 控制缓冲区大小

根据内存和性能需求调整 highWaterMark

7.4 流不是万能的

对于小文件(< 50MB),直接使用 fs.readFile 可能更简单。

总结

Node.js Stream 是处理大文件和构建数据管道的利器。通过 pipepipeline,你可以轻松组合多个流,实现高效的数据处理。本文介绍了流的基本用法、自定义转换流、背压处理以及实战案例。掌握这些技巧后,你将能够优雅地处理 GB 级别的数据。

延伸阅读

  • Node.js 官方文档:Stream API
  • 学习使用 stream.finished()stream.destroy() 进行资源清理
  • 探索第三方流库如 through2split2