人人都会AI编程

文件流式读写与大文件处理

更新时间:2026-07-11

在前几节中,我们已经掌握了 fs.readFilefs.writeFile 这类一次性将文件内容全部加载到内存的 API。然而,当文件体积达到几百 MB 甚至 GB 级别时,一次性读写不仅会占用海量内存,还可能导致服务假死直至内存溢出。文件流(Stream) 正是为解决这一问题而设计的,它将文件划分为小块(chunk),逐块进行处理,使得内存占用始终控制在一个可预测的范围内。

流式读写的核心概念

Node.js 的文件流基于 stream 模块构建,fs.createReadStreamfs.createWriteStream 返回的对象分别实现了可读流和可写流的接口。它们具备以下关键特性:

  • 分块传输:流内部维护一个缓冲区,每次从磁盘读取一小块数据,传递给下游,而不是等待整个文件读完。
  • 自动背压(Backpressure):当下游消费者处理速度慢于读取速度时,可读流会暂停继续读取,防止内存中堆积过多数据;待下游消费完毕,再恢复读取。
  • 事件驱动:通过 dataenderror 等事件,开发者可以精确控制数据处理流程,也可以使用 pipe 自动化连接各个流。

可读流:逐块读取大文件

使用 fs.createReadStream 创建一个可读流,再通过监听 data 事件逐块处理数据:

const fs = require('fs');

// 创建一个可读流,highWaterMark 指定每次读取的字节数(默认64KB)
const readStream = fs.createReadStream('./largefile.txt', {
  encoding: 'utf-8',
  highWaterMark: 16 * 1024 // 16KB per chunk
});

let lineCount = 0;
let leftover = '';

readStream.on('data', (chunk) => {
  // 每次处理一个 chunk,不会撑爆内存
  const lines = (leftover + chunk).split('\n');
  leftover = lines.pop(); // 保存不完整的行,与下一个chunk拼接
  lineCount += lines.length;
});

readStream.on('end', () => {
  if (leftover) lineCount++;
  console.log(`总行数:${lineCount}`);
});

readStream.on('error', (err) => {
  console.error('文件读取失败:', err.message);
});

在这个例子中,即使 largefile.txt 有 10GB,内存占用也始终只是 highWaterMark 左右的大小(加上少量 JavaScript 字符串开销)。leftover 变量的作用是处理跨 chunk 的断行情况,保证行数的统计结果准确。

可写流:逐块生成大文件

类似地,大文件的写入也应该使用 fs.createWriteStream,避免一次性将海量数据放在内存中再一起写入:

const writeStream = fs.createWriteStream('./output.txt', {
  encoding: 'utf-8',
  highWaterMark: 16 * 1024
});

for (let i = 0; i < 1000000; i++) {
  const line = `这是第 ${i} 行数据\n`;
  
  // write 方法返回 boolean,false 表示内部缓冲区已满,需要等待 drain 事件
  if (!writeStream.write(line)) {
    // 暂停写入,待缓冲区排空后再继续
    await new Promise(resolve => writeStream.once('drain', resolve));
  }
}

// 标记写入完成,并等待结束
writeStream.end(() => {
  console.log('大文件写入完成');
});

这里的关键在于 write() 的返回值:当内部缓冲区(highWaterMark)被填满时,write() 返回 false,此时应当等待 drain 事件再继续写入,否则数据会无限积压在内存中,失去流的意义。

管道(pipe):将读写串联的经典方式

最常用的流式操作是使用 pipe 方法,把可读流直接连接到可写流,整个流程自动处理背压和结束信号:

const readStream = fs.createReadStream('./source.tar.gz');
const writeStream = fs.createWriteStream('./backup.tar.gz');

readStream.pipe(writeStream);

如果需要中间加工,例如解压、加密,可以在管道中加入转换流(Transform Stream):

const { createGzip } = require('zlib'); // gzip 压缩转换流
const { pipeline } = require('stream/promises'); // Node.js 15+ 提供的 pipeline 方法

const readStream = fs.createReadStream('./data.txt');
const writeStream = fs.createWriteStream('./data.txt.gz');

// pipeline 会自动处理错误和流的关闭,比手动 pipe 更安全
await pipeline(
  readStream,
  createGzip(),
  writeStream
);
console.log('压缩完成');

pipeline 是优于 pipe 的推荐写法,因为它会在任何一个流发生错误时正确销毁所有流,并返回一个便于 async/await 的 Promise。

大文件处理的实战场景

1. 日志文件分析与过滤

面对上百 GB 的服务器访问日志,提取其中符合特定条件的条目并写入新文件:

const readStream = fs.createReadStream('./access.log');
const writeStream = fs.createWriteStream('./error.log');
const zlib = require('zlib');

readStream
  .pipe(require('split2')())           // 将文本流按行分割,split2 是第三方库
  .pipe(new require('stream').Transform({
    objectMode: true,
    transform(line, encoding, callback) {
      if (line.includes('500 Internal Server Error')) {
        // 只保留错误日志行,并写入压缩文件
        callback(null, line + '\n');
      } else {
        callback();
      }
    }
  }))
  .pipe(zlib.createGzip())
  .pipe(fs.createWriteStream('./error.log.gz'));

这个管道将日志按行分割、过滤、压缩、写入,处理过程中内存占用极低。

2. 分块上传/下载

在文件上传服务中,可以将客户端上传的分片直接流式写入磁盘,而无需等到所有分片收集完毕再合并:

// 接收上传流的简易示例
const http = require('http');
const fs = require('fs');

http.createServer((req, res) => {
  const writeStream = fs.createWriteStream('./uploaded_file');
  req.pipe(writeStream);  // 直接将 HTTP 请求的流管道到文件流

  req.on('end', () => {
    res.end('上传完成');
  });
});

req 本身是一个可读流,直接 pipe 到可写文件流,整个上传过程中的数据都是流式传递,不会在服务端积压。

3. 文件校验与转换

对文件边读边计算 MD5 或进行字符编码转换:

const crypto = require('crypto');

const readStream = fs.createReadStream('./raw_data');
const hash = crypto.createHash('sha256');
const transform = new require('stream').Transform({
  transform(chunk, encoding, callback) {
    // 同时传递 chunk 给 hash 和下游
    hash.update(chunk);
    this.push(chunk);
    callback();
  }
});

readStream
  .pipe(transform)
  .pipe(fs.createWriteStream('./data_copy'));

readStream.on('end', () => {
  console.log('SHA256:', hash.digest('hex'));
});

内存与性能的直观对比

用一个简单的测试来展示一次性读写和流式读写的差异。假设有一个 1GB 的文本文件:

一次性读取

// 内存将瞬间飙升至 ~1GB,甚至导致 OOM
const data = fs.readFileSync('./1gb.txt', 'utf-8');
fs.writeFileSync('./copy.txt', data);

流式复制

// 内存稳定在 ~20MB(取决于 highWaterMark)
fs.createReadStream('./1gb.txt').pipe(fs.createWriteStream('./copy.txt'));

在生产环境中,这种差异往往直接决定了服务在面对高并发大文件操作时能否稳定运行。

注意事项与常见问题

  1. 背压处理不当:自己实现 data 事件配合 write 时,必须检查 write 返回值和 drain 事件,否则流会失去背压保护。
  2. 错误处理:流中的错误不会自动传播,必须监听 error 事件,或者使用 pipeline 让错误冒泡为 Promise rejection。
  3. 对象模式与二进制模式:文件流默认是二进制(Buffer)模式,如果每个 chunk 需要按行或按对象处理,需要借助 split2JSONStream 等第三方转换流,并在构建 Transform 时设置 objectMode: true
  4. 高吞吐场景下的缓冲区大小highWaterMark 的默认值通常是 64KB,对于极高速的磁盘或网络,可以适当调大以减少系统调用次数,但会提高内存占用,需要根据实际环境压测调优。

小结

文件流式读写与管道机制是 Node.js 处理大文件的基石,也是“非阻塞 I/O”哲学在文件系统上的具体体现。掌握流的创建、背压控制、错误处理和常用组合模式,能够帮助我们在处理海量数据时编写出内存友好、性能健壮的代码。无论是日志处理、文件拷贝,还是网络上传下载,流都是不可或缺的利器。在后续章节中,我们还会进一步探索流的背压原理以及如何编写自定义的 Transform 流。