highWaterMark當讀取 read() 跟寫入 push(chunk) 的速度都是一致的
this.push("1".repeat(16384)); // 16384 bytes
myReadable.read(); // 16384 bytes
每次 push(chunk) 都會把 internal buffer 填滿,之後再用 read() 一次把 internal buffer 清空
但實務上讀取 read 跟寫入 push 的速度會不一樣,此時就需要記憶體管理機制,避免 OOM
push 的回傳值,如同 writable.write 一樣是 boolean,代表的是 isSafeToPushMore、canContinue 的意思
我們將 highWaterMark 設為 10,並且在 _read 使用 while 迴圈每次 push 6 bytes
import { Readable } from "stream";
class MyReadable extends Readable {
// 設定計數器,最多觸發 5 次 push(chunk)
private maxCount = 5;
private curCount = 0;
_read(size: number): void {
// 第 6 次就會 push(null) 來結束 readable
if (this.curCount === this.maxCount) {
this.push(null);
return;
}
// 模擬讀取資料的延遲
setTimeout(() => {
let isSafeToPushMore = true;
while (isSafeToPushMore) {
isSafeToPushMore = this.push("1".repeat(6));
console.log({
readableLength: this.readableLength,
isSafeToPushMore: isSafeToPushMore,
});
this.curCount++;
if (this.curCount === this.maxCount) break;
}
}, 100);
}
}
const myReadable = new MyReadable({ highWaterMark: 10 });
myReadable.on("readable", myReadable.read);
// Prints
// { readableLength: 6, isSafeToPushMore: true }
// { readableLength: 12, isSafeToPushMore: false }
// { readableLength: 6, isSafeToPushMore: true }
// { readableLength: 12, isSafeToPushMore: false }
// { readableLength: 6, isSafeToPushMore: true }
readableLength 代表目前 internal buffer 有多少 bytes 的資料等著被 read 讀取readableLength 總共 6 bytes,沒有頂到 highWaterMark,印出 { isSafeToPushMore: true }
readableLength 總共 12 bytes,頂到 highWaterMark,印出 { isSafeToPushMore: false }
{ isSafeToPushMore: false } 就不繼續寫入 internal bufferstream.Readable 只有三個 internal method,其中 _read 跟 _construct 的錯誤,最後都會傳遞到 _destroy
class MyReadable extends Readable {
_read(size: number): void;
_construct(callback: (error?: Error | null) => void): void;
_destroy(error: Error | null, callback: (error?: Error | null) => void): void;
}
_construct 階段呼叫 callback(err)import { Readable } from "stream";
class MyReadable extends Readable {
_construct(callback: (error?: Error | null) => void): void {
console.log("_construct");
callback(new Error("_construct failed"));
}
_destroy(
error: Error | null,
callback: (error?: Error | null) => void,
): void {
console.log("_destroy");
// ✅ _construct 拋出的錯誤會傳到 _destroy,請記得傳遞到 callback
callback(error);
}
}
const myReadable = new MyReadable();
// ❌ on("end") 不會觸發
myReadable.on("end", () => console.log('on("end")'));
// ✅ 使用者請記得用 on("error") 捕捉錯誤
myReadable.on("error", (err) => console.log('on("error")'));
myReadable.on("close", () => console.log('on("close")'));
執行順序
_construct;
_destroy;
on("error");
on("close");
_read 階段呼叫 destroy(err)import { Readable } from "stream";
import assert from "assert";
class MyReadable extends Readable {
_read(size: number): void {
console.log("_read");
this.destroy(new Error("_read failed"));
}
_destroy(
error: Error | null,
callback: (error?: Error | null) => void,
): void {
console.log("_destroy");
// ✅ destroy 背後會呼叫 _destroy,請記得把 error 傳遞到 callback
callback(error);
}
}
const myReadable = new MyReadable();
myReadable.on("readable", myReadable.read);
// ❌ on("end") 不會觸發
myReadable.on("end", () => console.log('on("end")'));
// ✅ 使用者請記得用 on("error") 捕捉錯誤
myReadable.on("error", (err) => console.log('on("error")'));
myReadable.on("close", () => console.log('on("close")'));
執行順序
_read;
_destroy;
on("error");
on("close");
writable._final vs readable.push在 stream.Readable,結束的訊號 readable.push(null) 是由實作者在 _read 的實作內主動呼叫的
_read(size: number): void {
// 實作者可以在這邊處理 async 操作
this.push(null);
}
在 stream.Writable,結束的訊號 writable.end() 是由使用者主動呼叫的
import { Writable } from "stream";
class MyWritable extends Writable {
_final(callback: (error?: Error | null) => void): void {
// 實作者可以在這邊處理 async 操作
}
}
const myWritable = new MyWritable();
myWritable.write("some data");
myWritable.end();
也因此,stream.Readable 沒有 _final 這個 internal method,因為 Node.js 已經提供在 _read 實作非同步操作的彈性了
在這篇文章,我們學到了
Readable 的記憶體管理Readable 的錯誤處理Readable 跟 Writable 的結束訊號差異Readable 跟 Writable