引言

Stream 是 Node.js 最核心的抽象之一。从文件读写、HTTP 请求响应到数据库连接,Stream 无处不在。然而,许多开发者对 Stream 的理解停留在"pipe 一下就好"的层面,对背压(backpressure)、错误传播、自定义流的实现知之甚少。

本文将深入 Stream 的内部机制,从四种基本流类型出发,结合实际场景,帮你真正掌握这一强大工具。


一、为什么需要 Stream?

传统方式的局限

假设你需要读取一个 2GB 的日志文件并写入到另一个位置:

javascript
// ❌ 传统方式:将整个文件读入内存
const fs = require('fs');

fs.readFile('/path/to/huge-file.log', (err, data) => {
  if (err) throw err;
  // 2GB 数据全部在内存中 —— 大多数服务器根本无法承受
  fs.writeFile('/path/to/output.log', data, (err) => {
    if (err) throw err;
  });
});

Stream 的解决方案

javascript
// ✅ Stream 方式:分块处理,内存占用恒定
const fs = require('fs');

const readStream = fs.createReadStream('/path/to/huge-file.log', {
  highWaterMark: 64 * 1024, // 每次读取 64KB
});

const writeStream = fs.createWriteStream('/path/to/output.log');

readStream.pipe(writeStream);

readStream.on('error', (err) => console.error('读取失败:', err));
writeStream.on('error', (err) => console.error('写入失败:', err));
writeStream.on('finish', () => console.log('文件复制完成'));

二、四种基本流类型

Node.js 提供了四种流抽象,覆盖了所有数据流动场景:

类型 方向 核心方法 典型场景
Readable 数据流出 _read() 文件读取、HTTP 请求体
Writable 数据流入 _write() 文件写入、HTTP 响应
Duplex 双向(独立) _read() + _write() TCP Socket、WebSocket
Transform 双向(转换) _transform() 压缩、加密、解析

Readable Stream

Readable 流负责生产数据。数据以两种模式流动:

javascript
const { Readable } = require('stream');

// 实现一个生成斐波那契数列的 Readable 流
class FibonacciStream extends Readable {
  constructor(max = 1000) {
    super({ objectMode: true });
    this.max = max;
    this.a = 0;
    this.b = 1;
  }

  _read() {
    const next = this.a + this.b;
    this.a = this.b;
    this.b = next;

    if (next > this.max) {
      this.push(null); // 结束流
      return;
    }

    this.push({ value: next, index: this.index++ });
  }
}

const fib = new FibonacciStream(10000);

fib.on('data', (chunk) => {
  console.log(chunk);
});

fib.on('end', () => {
  console.log('数列生成完毕');
});

Writable Stream

javascript
const { Writable } = require('stream');

// 实现一个批量写入数据库的 Writable 流
class BatchInsertStream extends Writable {
  constructor(db, options) {
    super({ objectMode: true, ...options });
    this.db = db;
    this.buffer = [];
    this.batchSize = 100;
  }

  async _write(chunk, encoding, callback) {
    this.buffer.push(chunk);

    if (this.buffer.length >= this.batchSize) {
      try {
        await this.db.insertMany(this.buffer);
        this.buffer = [];
        callback(); // 写入成功
      } catch (err) {
        callback(err); // 传播错误
      }
    } else {
      callback();
    }
  }

  async _final(callback) {
    // 处理剩余的缓冲区数据
    try {
      if (this.buffer.length > 0) {
        await this.db.insertMany(this.buffer);
      }
      callback();
    } catch (err) {
      callback(err);
    }
  }
}

Transform Stream

Transform 流是最常用也最强大的流类型,同时实现了可读和可写接口:

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

// JSON 行解析 Transform 流
class JSONLineParser extends Transform {
  constructor(options) {
    super({ readableObjectMode: true, ...options });
    this.buffer = '';
  }

  _transform(chunk, encoding, callback) {
    this.buffer += chunk.toString();

    // 按行分割并解析 JSON
    const lines = this.buffer.split('\n');
    // 最后一行可能不完整,保留在下一次处理
    this.buffer = lines.pop() || '';

    for (const line of lines) {
      if (line.trim()) {
        try {
          this.push(JSON.parse(line));
        } catch (err) {
          callback(err);
          return;
        }
      }
    }
    callback();
  }

  _flush(callback) {
    // 处理最终的缓冲区
    if (this.buffer.trim()) {
      try {
        this.push(JSON.parse(this.buffer));
      } catch (err) {
        callback(err);
        return;
      }
    }
    callback();
  }
}

三、背压机制深度解析

背压(backpressure)是 Stream 最容易被忽视但最关键的概念。当数据生产速度超过消费速度时,如果放任不管,内存会被无限增长的内部缓冲区撑爆。

背压的工作原理

javascript
const fs = require('fs');

const readable = fs.createReadStream('source.dat', {
  highWaterMark: 16 * 1024, // 16KB 内部缓冲上限
});

const writable = fs.createWriteStream('dest.dat');

// readable.pipe(writable) 内部的实际流程:
readable.on('data', (chunk) => {
  // writable.write() 返回 false 表示内部缓冲区已满
  const canContinue = writable.write(chunk);

  if (!canContinue) {
    // 暂停读取,等待 writable 排空
    readable.pause();

    // 当 writable 排空时恢复读取
    writable.once('drain', () => {
      readable.resume();
    });
  }
});

writable.on('drain', () => {
  console.log('缓冲区排空,ready 继续写入');
});

pipe() 的完整手动实现

理解 pipe() 的实现有助于掌握背压:

javascript
function manualPipe(readable, writable) {
  // 错误处理:任一流出错,销毁另一方
  const onReadableError = (err) => {
    writable.destroy(err);
  };
  const onWritableError = (err) => {
    readable.destroy(err);
  };

  readable.on('error', onReadableError);
  writable.on('error', onWritableError);

  readable.on('data', (chunk) => {
    const canWrite = writable.write(chunk);

    if (!canWrite) {
      readable.pause();

      // writable 排空后恢复 readable
      writable.once('drain', () => {
        readable.resume();
      });
    }
  });

  // readable 结束时关闭 writable
  readable.on('end', () => {
    writable.end();
  });

  return writable;
}

highWaterMark 调优

javascript
// 不同场景下的 highWaterMark 选择
const fs = require('fs');

// 大文件复制:较大的缓冲区减少 I/O 调用次数
const bigFileRead = fs.createReadStream('large-video.mp4', {
  highWaterMark: 1024 * 1024, // 1MB
});

// 实时日志处理:较小的缓冲区降低延迟
const logRead = fs.createReadStream('app.log', {
  highWaterMark: 4 * 1024, // 4KB
});

// 对象模式:highWaterMark 是对象数量而非字节
const { Readable } = require('stream');
const objStream = new Readable({
  objectMode: true,
  highWaterMark: 16, // 16 个对象
  read() {},
});

四、pipeline:现代 Stream 编排

stream.pipeline 是 pipe() 的现代替代品,提供了统一的错误处理:

javascript
const { pipeline } = require('stream/promises');
const fs = require('fs');
const zlib = require('zlib');
const crypto = require('crypto');

async function compressAndEncrypt(inputPath, outputPath, password) {
  try {
    await pipeline(
      fs.createReadStream(inputPath),
      zlib.createGzip(),
      crypto.createCipheriv('aes-256-gcm', key, iv),
      fs.createWriteStream(outputPath),
    );
    console.log('压缩加密完成');
  } catch (err) {
    // pipeline 会自动销毁所有流并传播错误
    console.error('管道处理失败:', err);
  }
}

pipeline vs pipe 对比

特性 pipe() pipeline()
错误传播 不自动传播 自动传播到回调/Promise
流清理 需手动 destroy 自动销毁所有中间流
回调/Promise 仅 EventEmitter 支持 callback 和 Promise
可组合性 链式调用 数组形式,可读性更好

五、实战案例

案例 1:CSV 大文件解析与数据库导入

javascript
const fs = require('fs');
const { pipeline } = require('stream/promises');
const { Transform } = require('stream');

// CSV 行解析器
class CSVParser extends Transform {
  constructor(headers, options) {
    super({ readableObjectMode: true, ...options });
    this.headers = headers;
    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 values = line.split(',');
      const record = {};
      this.headers.forEach((header, i) => {
        record[header] = values[i]?.trim() ?? '';
      });
      this.push(record);
    }
    callback();
  }

  _flush(callback) {
    if (this.buffer.trim()) {
      const values = this.buffer.split(',');
      const record = {};
      this.headers.forEach((header, i) => {
        record[header] = values[i]?.trim() ?? '';
      });
      this.push(record);
    }
    callback();
  }
}

// 批量插入 Transform
class BatchInsertTransformer extends Transform {
  constructor(db, batchSize = 1000) {
    super({ objectMode: true });
    this.db = db;
    this.batchSize = batchSize;
    this.batch = [];
  }

  async _transform(record, encoding, callback) {
    this.batch.push(record);

    if (this.batch.length >= this.batchSize) {
      try {
        await this.db.insertMany(this.batch);
        console.log(`已插入 ${this.batch.length} 条记录`);
        this.batch = [];
        callback();
      } catch (err) {
        callback(err);
      }
    } else {
      callback();
    }
  }

  async _flush(callback) {
    try {
      if (this.batch.length > 0) {
        await this.db.insertMany(this.batch);
        console.log(`已插入最后 ${this.batch.length} 条记录`);
      }
      callback();
    } catch (err) {
      callback(err);
    }
  }
}

// 使用
async function importCSV(filePath, db) {
  const headers = ['id', 'name', 'email', 'created_at', 'status'];

  await pipeline(
    fs.createReadStream(filePath, { highWaterMark: 64 * 1024 }),
    new CSVParser(headers),
    new BatchInsertTransformer(db, 500),
  );

  console.log('导入完成');
}

案例 2:实时日志过滤与转发

javascript
const { pipeline } = require('stream/promises');
const { Transform } = require('stream');
const net = require('net');

// 日志级别过滤器
class LogFilter extends Transform {
  constructor(minLevel, options) {
    super({ ...options });
    this.minLevel = minLevel;
  }

  _transform(chunk, encoding, callback) {
    const line = chunk.toString().trim();

    // 检查日志级别
    const levelMatch = line.match(/\[(ERROR|WARN|INFO|DEBUG)\]/);

    if (levelMatch) {
      const level = levelMatch[1];
      const shouldForward = this.shouldLog(level);

      if (shouldForward) {
        this.push(chunk);
      }
    }
    callback();
  }

  shouldLog(level) {
    const levels = { DEBUG: 0, INFO: 1, WARN: 2, ERROR: 3 };
    return levels[level] >= levels[this.minLevel];
  }
}

// 使用标准输入作为日志源
async function filterStdin() {
  await pipeline(
    process.stdin,
    new LogFilter('WARN'),
    process.stdout,
  );
}

// 命令行使用:
// tail -f /var/log/app.log | node log-filter.js

案例 3:HTTP 大文件下载(支持断点续传)

javascript
const http = require('http');
const fs = require('fs');
const { pipeline } = require('stream/promises');

async function downloadFile(url, destPath) {
  return new Promise((resolve, reject) => {
    http.get(url, (response) => {
      // 处理重定向
      if (response.statusCode >= 300 && response.statusCode < 400) {
        return downloadFile(response.headers.location, destPath)
          .then(resolve)
          .catch(reject);
      }

      if (response.statusCode !== 200) {
        reject(new Error(`下载失败: HTTP ${response.statusCode}`));
        return;
      }

      const totalSize = parseInt(response.headers['content-length'], 10);
      let downloadedSize = 0;

      // 使用 Transform 流监控进度
      const progressTracker = new Transform({
        transform(chunk, encoding, callback) {
          downloadedSize += chunk.length;
          const progress = ((downloadedSize / totalSize) * 100).toFixed(1);
          process.stdout.write(`\r下载进度: ${progress}%`);
          this.push(chunk);
          callback();
        },
      });

      pipeline(
        response,
        progressTracker,
        fs.createWriteStream(destPath),
      )
        .then(() => {
          console.log('\n下载完成');
          resolve();
        })
        .catch(reject);
    }).on('error', reject);
  });
}

六、性能与内存优化

内存使用对比

javascript
const fs = require('fs');

// 测试:处理 1GB 文件的不同方式内存占用

// 方式 A:readFile —— 内存 ≈ 1GB + overhead
// fs.readFile('1gb-file.dat', callback);

// 方式 B:Stream —— 内存 ≈ 64KB(highWaterMark)
// fs.createReadStream('1gb-file.dat', { highWaterMark: 64 * 1024 });

// 方式 C:Stream + 处理 —— 内存 ≈ 64KB + 处理开销
// stream.pipe(transform).pipe(writeStream);

并发流控制

javascript
const { pipeline } = require('stream/promises');
const { Transform } = require('stream');

// 并发处理 Transform:允许 N 个异步任务同时进行
class ParallelTransform extends Transform {
  constructor(concurrency, transformFn, options) {
    super({ objectMode: true, ...options });
    this.concurrency = concurrency;
    this.transformFn = transformFn;
    this.active = 0;
    this.queue = [];
  }

  _transform(chunk, encoding, callback) {
    this.active++;

    this.transformFn(chunk)
      .then((result) => {
        this.push(result);
        this.active--;
        this._processQueue();
      })
      .catch(callback);

    if (this.active < this.concurrency) {
      callback();
    } else {
      this.queue.push(callback);
    }
  }

  _processQueue() {
    while (this.active < this.concurrency && this.queue.length > 0) {
      const cb = this.queue.shift();
      cb();
    }
  }
}

// 使用:并发 5 个请求批量处理记录
async function processWithConcurrency() {
  const transform = new ParallelTransform(5, async (record) => {
    const result = await fetch(`https://api.example.com/process`, {
      method: 'POST',
      body: JSON.stringify(record),
    });
    return result.json();
  });

  await pipeline(
    sourceStream,
    transform,
    destinationStream,
  );
}

七、常见陷阱与最佳实践

1. 忘记处理 error 事件

javascript
// ❌ 未处理的 error 会导致进程崩溃
readStream.pipe(writeStream);

// ✅ 始终处理错误
readStream.on('error', handleError);
writeStream.on('error', handleError);

// ✅ 使用 pipeline 自动传播错误
await pipeline(readStream, writeStream);

2. 混合使用 pipe 和事件

javascript
// ❌ 同时使用 pipe 和 data 事件
readStream.pipe(writeStream);
readStream.on('data', (chunk) => {
  // 数据会被两个消费者竞争
});

// ✅ 需要观测数据时使用 Transform 流
const observer = new Transform({
  transform(chunk, encoding, callback) {
    console.log(`处理了 ${chunk.length} 字节`);
    this.push(chunk);
    callback();
  },
});

readStream.pipe(observer).pipe(writeStream);

3. 流完成后的清理

javascript
async function processFile(input, output) {
  const readStream = fs.createReadStream(input);
  const writeStream = fs.createWriteStream(output);

  try {
    await pipeline(readStream, writeStream);
    // 流已自动关闭,无需额外清理
  } catch (err) {
    // pipeline 失败时自动销毁所有流
    console.error('处理失败:', err);
  }
}

结语

Stream 是 Node.js 的基石——它不仅是 API,更是一种编程范式。掌握 Stream 意味着你可以在处理数据时保持恒定的内存占用、天然的背压支持和优雅的组合能力。

核心要点:

  1. 四种流类型对应不同的数据流向需求,理解 _read、_write、_transform 是关键
  2. 背压是 Stream 设计中最精妙的部分,highWaterMark 和 drain 事件是核心机制
  3. pipeline 取代 pipe() 成为现代准则,提供自动错误传播和流清理
  4. 自定义流让你可以像搭积木一样组合数据处理管道
  5. 始终处理错误——未捕获的 Stream error 会直接崩溃 Node.js 进程

当你下次面对大文件、实时数据或高并发场景时,首先问自己:这能不能用 Stream 解决?