ประมวลผลไฟล์ขนาดหลายกิกะไบต์โดยใช้หน่วยความจำคงที่ ด้วยสตรีม pipeline และ backpressure
เปิดบทเรียนนี้ใน KodokonreadFile โหลดไฟล์ทั้งไฟล์เข้าสู่หน่วยความจำ บนล็อกขนาด 4 GB โปรเซสของคุณจะระเบิด สตรีม (stream) ประมวลผลข้อมูลเป็นชิ้น ๆ (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() ทำได้ดียิ่งกว่า นอกจาก backpressure แล้ว มันยังส่งต่อข้อผิดพลาดและทำลายสตรีมทั้งหมดอย่างเรียบร้อยในทั้งสองทิศทาง
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 เป็น async iterable นั่นคือ for await จะบริโภคชิ้นข้อมูลไปพร้อมกับเคารพ backpressure ระวังไว้ ชิ้นข้อมูล (chunk) คือการตัดไบต์ออกมาแบบ ตามอำเภอใจ ไม่มีอะไรรับประกันว่ามันจะจบลงตรงขอบบรรทัดพอดี สำหรับการประมวลผลทีละบรรทัด node:readline จะจัดการประกอบชิ้นส่วนกลับเข้าด้วยกัน รวมถึงกรณี \r\n ที่ถูกแบ่งข้ามสองชิ้น ด้วย crlfDelay: Infinity
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");