Kodokon kodokon.com

流与背压:处理大文件

借助流、pipeline 和背压,以恒定内存处理多 GB 的文件。

9 分钟 · 3 题

在 Kodokon 中打开本课

readFile 会把整个文件加载进内存:面对一个 4 GB 的日志,你的进程会崩溃。流(Streams)以块(chunk)为单位处理数据,占用恒定的内存,其上限由 highWaterMark 决定(文件流默认为 64 KiB)。共有四种类型:ReadableWritableDuplexTransform。它们解决的核心问题被称为背压(backpressure):当生产者比消费者更快时,你该怎么办?

JAVASCRIPT
import { createWriteStream } from "node:fs";

const out = createWriteStream("big.txt");
let i = 0;

function writeChunks() {
  let ok = true;
  while (i < 1e6 && ok) {
    ok = out.write(`line ${i}\n`);
    i += 1;
  }
  if (i < 1e6) out.once("drain", writeChunks);
  else out.end();
}

writeChunks();
手工实现的 write/drain 契约。

write() 返回 false 时,内部缓冲区已经超过了 highWaterMark:数据并没有丢失,但继续写入会把所有内容都堆积在内存中 - 这恰恰是流本应避免的。契约是:停止写入并等待 drain 事件。这正是 pipe() 为你所做的,而 pipeline() 做得更好:在背压之上,它还会传播错误并干净地销毁双向的所有流。

JAVASCRIPT
import { createReadStream, createWriteStream }
  from "node:fs";
import { Transform } from "node:stream";
import { pipeline } from "node:stream/promises";

const upper = new Transform({
  transform(chunk, encoding, callback) {
    callback(null, chunk.toString().toUpperCase());
  },
});

await pipeline(
  createReadStream("big.txt"),
  upper,
  createWriteStream("big-upper.txt"),
);
console.log("done");
流式转换,恒定内存。

Readable 流是异步可迭代对象:for await 会在遵守背压的前提下逐块消费数据。要小心,一个 chunk 是一段任意的字节切片:没有任何东西能保证它正好在一个行边界处结束。对于逐行处理,node:readline 会负责把这些碎片重新拼接起来,包括借助 crlfDelay: Infinity 处理跨两个 chunk 分开的 \r\n

JAVASCRIPT
import { createReadStream } from "node:fs";
import { createInterface } from "node:readline";

const rl = createInterface({
  input: createReadStream("big.txt"),
  crlfDelay: Infinity,
});

let count = 0;
for await (const line of rl) {
  if (line.includes("42")) count += 1;
}
console.log(count, "matching lines");
逐行读取一个大文件。

知识检测

确认你已牢记本课的重点内容。

  1. 当 writable.write(chunk) 返回 false 时意味着什么?
    • chunk 丢失了,必须重新写入
    • 内部缓冲区超过了 highWaterMark:停止写入直到 drain 事件
    • 流被一个错误销毁了
  2. pipeline() 相较于 pipe() 有什么决定性的优势?
    • pipeline 更快,因为它是多线程的
    • pipe 处理背压,pipeline 不处理
    • pipeline 会传播错误并销毁所有涉及的流
    • pipeline 会自动转换编码
  3. 在 objectMode 下,highWaterMark 代表什么?
    • 队列中对象的数量
    • 字节数,就像在二进制模式下那样
    • 序列化后对象的最大大小