ข้ามไปยังเนื้อหา

Backpressure

ลองนึกภาพสายยักษ์ต่อเข้ากับหัวฉีดน้ำในสวน ถ้าน้ำไหลเข้าเร็วกว่าที่หัวฉีดจะพ่นออกได้ สายก็จะพองและแตกในที่สุด เหตุการณ์เดียวกันนี้เกิดขึ้นกับ Node Stream: ถ้า Readable source ผลิตข้อมูลเร็วกว่าที่ Writable จะรับได้ chunk ที่ Buffer สะสมใน Memory จนกระบวนการหมด RAM หรือ crash

Backpressure คือกลไกที่ consumer ที่ช้าบอกให้ producer ที่เร็วชะลอลง เทียบได้กับ flow control ใน TCP

Method write() ของ Writable Stream เป็นสัญญาณสำคัญ:

  • คืนค่า true — internal buffer ต่ำกว่า high-water mark (ค่าเริ่มต้น: 16 KB (16384 bytes) สำหรับ byte stream หรือ 16 objects สำหรับ object-mode stream โดย fs.createReadStream เป็นข้อยกเว้นที่ใช้ 64 KB) producer ส่งข้อมูลต่อได้
  • คืนค่า false — buffer อยู่ที่หรือเกิน high-water mark producer ต้อง pause และรอ event drain ก่อนเขียนต่อ

high-water mark ไม่ใช่ขีดจำกัดแบบ hard — เป็นเพียง soft advisory Node ยัง Buffer chunk เกินขีดนั้นได้แต่ต้องเสียค่า Memory การเคารพสัญญาณนี้ทำให้การใช้ Memory คงที่

ก่อนที่จะมี pipe และ pipeline คุณต้องจัดการเองดังนี้:

function pumpWithBackpressure(readable, writable) {
readable.on('data', function(chunk) {
var ok = writable.write(chunk);
if (!ok) {
readable.pause(); // หยุด source
writable.once('drain', function() {
readable.resume(); // resume เมื่อพร้อม
});
}
});
readable.on('end', function() {
writable.end();
});
}

Pattern นี้ verbose และเกิดข้อผิดพลาดได้ง่าย ลืม readable.pause() ครั้งเดียว buffer ก็เติบโตไม่หยุด

readable.pipe(writable) และ stream.pipeline() implement pattern ข้างต้นภายใน:

  1. ส่ง chunk จาก readable ไปยัง writable
  2. Pause readable เมื่อ write() คืนค่า false
  3. Resume readable เมื่อเกิด event drain

คุณไม่จำเป็นต้องจัดการการ pause และ resume เองเมื่อใช้ pipe หรือ pipeline นี่คือเหตุผลหลักที่ควรใช้แทน raw event listener

Demo ด้านล่างสร้าง Readable และ Writable ที่มี high-water mark เล็กมาก (1 byte) เพื่อให้ backpressure trigger ทุก chunk สังเกตว่า readable ถูก pause และ resume โดยอัตโนมัติ

Node.js

Needs the Node.js runtime — open in StackBlitz to run.

writable.write() คืนค่า false หมายความว่าอะไร?
ควร listen event ใดเพื่อรู้ว่าปลอดภัยที่จะ resume การเขียนหลัง backpressure?
readable.pipe() และ stream.pipeline() ช่วยเรื่อง backpressure อย่างไร?