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

Pipe & Pipeline

.pipe() เชื่อม Readable เข้ากับ Writable โดยตรงและจัดการการไหลของข้อมูลระหว่างกันโดยอัตโนมัติ รองรับ backpressure: เมื่อ writable บอกว่า buffer เต็ม (โดยคืนค่า false จาก write()) pipe จะ pause readable และ resume เมื่อ drain

const fs = require('node:fs');
const readable = fs.createReadStream('./input.txt');
const writable = fs.createWriteStream('./output.txt');
readable.pipe(writable);

คุณสามารถ chain .pipe() ผ่าน Transform Stream เพื่อสร้าง processing pipeline:

const zlib = require('node:zlib');
fs.createReadStream('./file.txt')
.pipe(zlib.createGzip())
.pipe(fs.createWriteStream('./file.txt.gz'));

.pipe() มีข้อบกพร่องที่ทราบกันดี: ไม่ส่งต่อ error ถ้า Stream ใดใน chain emit event error Stream อื่นจะไม่ถูกทำลายหรือปิดโดยอัตโนมัติ ส่งผลให้มี file handle ที่ยังเปิดอยู่และ memory leak

// อันตราย: error บน readable ไม่ปิด writable
readable.pipe(writable);
readable.on('error', function(err) {
console.error('read error — but writable is still open!', err.message);
writable.destroy(); // ต้อง clean up เอง
});

stream.pipeline(...streams, callback) ถูกเพิ่มใน Node 10 เพื่อแก้ปัญหานี้โดยเฉพาะ:

  1. เชื่อม Stream ทั้งหมดแบบ end-to-end
  2. ส่งต่อ error จาก Stream ใดก็ได้ไปยัง callback
  3. ทำลาย Stream ทุกตัวใน chain โดยอัตโนมัติเมื่อเกิด error
const { pipeline } = require('node:stream');
const fs = require('node:fs');
const zlib = require('node:zlib');
pipeline(
fs.createReadStream('./input.txt'),
zlib.createGzip(),
fs.createWriteStream('./input.txt.gz'),
function(err) {
if (err) {
console.error('pipeline failed:', err.message);
} else {
console.log('pipeline complete');
}
}
);

Node 15+ ยังมี stream/promises ที่ให้ใช้กับ await ได้:

const { pipeline } = require('node:stream/promises');
await pipeline(
fs.createReadStream('./input.txt'),
zlib.createGzip(),
fs.createWriteStream('./input.txt.gz')
);

Demo ด้านล่างสร้าง in-memory pipeline สั้นๆ: Readable source ส่งสาม string, Transform แปลงแต่ละ chunk เป็นตัวพิมพ์ใหญ่, และ Writable รวบรวมผลลัพธ์

Node.js

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

ปัญหาหลักของการใช้ readable.pipe(writable) เพียงอย่างเดียวคืออะไร?
เกิดอะไรขึ้นเมื่อ Stream ใดใน stream.pipeline() chain emit error?
Stream ชนิดใดที่วางกลาง pipeline เพื่อแปลงข้อมูลระหว่างทางได้?