Skip to content

Node系列 · Node基础:文件流

大文件(GB 级别)一次性 readFile 会撑爆内存;网络请求的 body 也是按块到达的——这两种"数据分批到达"的场景都靠流处理。理解流的四种类型和背压机制,就理解了 Node I/O 的精髓。

一、什么是流

流(Stream) 是一种按时间分批处理数据的抽象接口:

  • 数据不是一次性到达,而是分块(chunk)逐次可用
  • 内存占用只与"单 chunk 大小"有关,与"总数据量"无关
  • 生产者(写入端)和消费者(读取端)通过事件解耦

为什么需要流

场景readFile用流
10KB 配置✅ 简单❌ 杀鸡用牛刀
500MB 日志❌ 一次性占用 500MB 内存✅ 始终 64KB 缓冲区
1GB 视频❌ 直接 OOM✅ 边读边处理
网络大文件下载❌ 全部缓冲在内存✅ 流式写盘

二、四种流类型

Node 内置 4 种抽象流:

类型抽象类方向示例
Readablestream.Readable只读fs.createReadStreamhttp.IncomingMessage
Writablestream.Writable只写fs.createWriteStreamhttp.ServerResponse
Duplexstream.Duplex读写net.Socketzlib.createGzip
Transformstream.Transform读写 + 转换zlib.createGzipcrypto.createCipheriv

三、文件流

3.1 创建文件流

javascript
const fs = require('node:fs');

const reader = fs.createReadStream('./big.txt', {
  encoding: 'utf-8',     // 不指定则 chunk 是 Buffer
  highWaterMark: 64 * 1024, // 默认 64KB;可调
});

const writer = fs.createWriteStream('./out.txt', {
  encoding: 'utf-8',
  flags: 'w',             // 默认 'w',会覆盖
});

highWaterMark 是"水位线"——内部缓冲超过这个值就触发背压。

3.2 消费流:事件 vs async iterator

事件方式(经典)

javascript
const reader = fs.createReadStream('./big.txt', { encoding: 'utf-8' });

reader.on('data', (chunk) => {
  console.log(`收到 ${chunk.length} 字节`);
});

reader.on('end', () => {
  console.log('读取完成');
});

reader.on('error', (err) => {
  console.error('读取失败:', err);
});

async iterator(推荐)

javascript
const fs = require('node:fs');

async function countLines(filePath) {
  const reader = fs.createReadStream(filePath, { encoding: 'utf-8' });
  let count = 0;
  for await (const chunk of reader) {
    count += chunk.split('\n').length - 1;
  }
  return count;
}

console.log(await countLines('./access.log'));

for await...of 自动处理 data / end / error 事件,代码更直观。

3.3 写入流

javascript
const writer = fs.createWriteStream('./out.txt');

writer.write('第一行\n');
writer.write('第二行\n');

writer.end();  // 结束写入;end 之后再 write 会报错

writer.on('finish', () => {
  console.log('全部写入完成');
});

writer.on('error', (err) => {
  console.error('写入失败:', err);
});

四、pipe:流的串联

pipe 把 readable 的输出直接灌进 writable,省去手动监听 data 事件:

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

await pipeline(
  fs.createReadStream('./source.txt'),
  fs.createWriteStream('./dest.txt')
);
console.log('复制完成');

pipeline(推荐用 stream/promises 版本)相比 reader.pipe(writer) 的优势:

方案错误传播资源清理
reader.pipe(writer)错误需各自监听失败时不会自动关另一端
pipeline(reader, writer)任一环节错误都会 reject失败时自动关闭所有流

多级 pipe 组合

javascript
await pipeline(
  fs.createReadStream('./access.log'),
  new Transform({
    transform(chunk, encoding, callback) {
      // 只保留包含 ERROR 的行
      const lines = chunk.toString().split('\n').filter((l) => l.includes('ERROR'));
      callback(null, lines.join('\n'));
    },
  }),
  fs.createWriteStream('./errors.log'),
);

五、背压(Backpressure)

背压是流处理最关键的概念。当写入端消化速度跟不上读取端产生速度,会发生:

读取端不断产生 chunk

写入端内部缓冲不断增长

超过 highWaterMark 水位

内存膨胀 → 可能 OOM

pipe 会自动处理背压:write() 返回 false 时暂停 readabledrain 事件触发后再恢复。手动实现就要自己监听:

javascript
const writer = fs.createWriteStream('./out.txt');

function write(data) {
  if (!writer.write(data)) {
    // 缓冲已满,等 drain 再写
    return new Promise((resolve) => writer.once('drain', resolve));
  }
  // 缓冲还有空间,立即返回
  return Promise.resolve();
}

await write('chunk 1\n');
await write('chunk 2\n');
await new Promise((resolve) => writer.end(resolve));

WARNING

不要忽略 write() 的返回值。直接 writer.write(chunk) 不管返回值,等同于背压机制不存在——最终会撑爆内存或触发 EPIPE 错误。

六、对象模式(Object Mode)

默认流传输 Buffer / string。可以开启 object mode 传输任意 JS 值:

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

const csv = new Transform({
  readableObjectMode: true,
  writableObjectMode: true,
  transform(row, encoding, callback) {
    // 输入是 JS 对象,输出也是 JS 对象
    callback(null, { ...row, processedAt: Date.now() });
  },
});

适用场景:

  • 逐行解析 CSV,每行是对象
  • 数据库查询流式消费,每条记录是对象
  • JSON 数组按对象迭代

七、流的事件完整对照

事件触发时机流类型
data每个 chunk 可读时Readable
end数据读完Readable
error出错(必须监听,否则崩溃)全部
drain写入缓冲降到水位以下,可继续写Writable
finish全部数据已写入底层Writable
close流及其底层资源已关闭全部
pipe / unpipe被 / 被取消 pipeReadable
pause / resume暂停 / 恢复推送Readable

八、流的两种工作模式

模式默认行为
流动模式(flowing)data 事件被监听时自动进入数据尽快推给消费者
暂停模式(paused)默认数据暂存内部缓冲,等显式 read()resume()

for await...of 内部会自动切换两种模式,开发者无感。

九、最佳实践

场景推荐
大文件处理(GB 级)必用流
HTTP body 接收已经是流(req 是 Readable)
HTTP body 发送已经是流(res 是 Writable)
多级流(解压 → 解密 → 解析)pipeline 串接
错误处理stream/promisespipeline / finished Promise 版本
背压pipe / pipeline 自动处理;手动实现要监听 drain
不要做的await reader.text() 把整个流读成字符串 / JSON.parse(await reader.text())

十、小结

  • 流的核心价值:内存占用与数据总量无关
  • 4 种类型:Readable(读)/ Writable(写)/ Duplex(读写)/ Transform(读写 + 转换)
  • 消费方式:on('data') 事件 或 for await...of(推荐)
  • 多级流用 pipeline 串联,自动处理错误传播和资源清理
  • 背压是流处理的灵魂——pipe / pipeline 自动处理,手写必须监听 drain
  • 文件系统流:fs.createReadStream / fs.createWriteStream,默认 highWaterMark 64KB