借助流、pipeline 和背压,以恒定内存处理多 GB 的文件。
在 Kodokon 中打开本课readFile 会把整个文件加载进内存:面对一个 4 GB 的日志,你的进程会崩溃。流(Streams)以块(chunk)为单位处理数据,占用恒定的内存,其上限由 highWaterMark 决定(文件流默认为 64 KiB)。共有四种类型:Readable、Writable、Duplex 和 Transform。它们解决的核心问题被称为背压(backpressure):当生产者比消费者更快时,你该怎么办?
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() 返回 false 时,内部缓冲区已经超过了 highWaterMark:数据并没有丢失,但继续写入会把所有内容都堆积在内存中 - 这恰恰是流本应避免的。契约是:停止写入并等待 drain 事件。这正是 pipe() 为你所做的,而 pipeline() 做得更好:在背压之上,它还会传播错误并干净地销毁双向的所有流。
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。
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");