Kodokon kodokon.com

สตรีมและ backpressure: การจัดการไฟล์ขนาดใหญ่

ประมวลผลไฟล์ขนาดหลายกิกะไบต์โดยใช้หน่วยความจำคงที่ ด้วยสตรีม pipeline และ backpressure

9 นาที · 3 คำถาม

เปิดบทเรียนนี้ใน Kodokon

readFile โหลดไฟล์ทั้งไฟล์เข้าสู่หน่วยความจำ บนล็อกขนาด 4 GB โปรเซสของคุณจะระเบิด สตรีม (stream) ประมวลผลข้อมูลเป็นชิ้น ๆ (chunk) โดยใช้หน่วยความจำคงที่ ซึ่งถูกจำกัดด้วย highWaterMark (ค่าเริ่มต้น 64 KiB สำหรับสตรีมไฟล์) มีอยู่สี่ชนิด ได้แก่ Readable, Writable, Duplex และ Transform ปัญหาหลักที่พวกมันแก้เรียกว่า 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() ทำได้ดียิ่งกว่า นอกจาก backpressure แล้ว มันยังส่งต่อข้อผิดพลาดและทำลายสตรีมทั้งหมดอย่างเรียบร้อยในทั้งสองทิศทาง

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

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 หมายความว่าอย่างไร?
    • ชิ้นข้อมูลสูญหายและต้องเขียนใหม่
    • บัฟเฟอร์ภายในเกิน highWaterMark ให้หยุดเขียนจนกว่าจะเกิดอีเวนต์ drain
    • สตรีมถูกทำลายเพราะข้อผิดพลาด
  2. pipeline() มีข้อได้เปรียบชี้ขาดอะไรเหนือ pipe()?
    • pipeline เร็วกว่าเพราะเป็นแบบหลายเธรด
    • pipe จัดการ backpressure ส่วน pipeline ไม่
    • pipeline ส่งต่อข้อผิดพลาดและทำลายสตรีมทั้งหมดที่เกี่ยวข้อง
    • pipeline แปลงการเข้ารหัสให้โดยอัตโนมัติ
  3. ใน objectMode นั้น highWaterMark หมายถึงอะไร?
    • จำนวนอ็อบเจกต์ที่อยู่ในคิว
    • จำนวนไบต์ เหมือนในโหมดไบนารี
    • ขนาดสูงสุดของอ็อบเจกต์ที่ถูก serialize