積み上がるバイトを自分で測る
目標
速い生産者と遅い消費者の間でメモリがどう増えていくかを自分で測り、write()の戻り値とdrainイベントが何を防いでくれるかを確認します。
なぜ重要なのか
「メモリリーク」として報告されるものの多くは、漏れているのではなく行列ができているのです。毎秒10MBを作り出す側と毎秒2MBを受け取る側をそのままつなぐと、余る8MBはどこかに溜まるしかなく、そのどこかがこのプロセスのヒープです。
Nodeのストリームは、これを防ぐためのつまみを1つ提供します。write()の戻り値です。ところがその値は、放っておくと何もしません。読んで止まるのは書く側の役目です。 その1行が抜けてもコードは問題なく動き、テストも通り、本番に出て大きなファイル1つでプロセスが死にます。
そして背圧を守ると遅くなりそうですが、測るとそうではありません。消費者が受け取る速度はどのみち同じだからです。失うものはほとんどなく、防げるものは大きいです。
ステップ
/root/work/backpressure/flow.mjsにmakeSink(options)を作り、 このバージョンのデフォルトのhighWaterMarkを測ってreport.jsonのhwmに書きます。fillUntilFalse(stream, chunk)でバッファーが満ちる地点を数えます。writeAll(stream, chunks)で背圧を守りながら書きます。floodNoWait(stream, chunks)で戻り値を無視して流し込み、runs.floodに書きます。- 同じ量を
writeAllで送ってruns.pacedに書き、backpressureに比率を書きます。 pipeThrough(readable, writable)で同じことをして、runs.pipelineに書きます。estimateQueueBytesとsecondsUntilで溜まる速度を計算します。
参考
makeSink({hwm, delayMs})は、受け取った量をreceivedBytes・receivedChunksに 数えるWritableです。delayMsが0ならsetImmediateで、0より大きければ その分だけ遅れてコールバックを呼びます。writeAllとfloodNoWaitは、どちらも{peakBuffered}を返します。stream.writableLengthがその瞬間に溜まっているバイト数です。- ステップ4とステップ5は同じ量で測らないと比較になりません。断片を配列であらかじめ作らず、 送るたびに新しく作ってください。あらかじめ作ると測る前にすでに使い切った状態になり、 同じバッファーを再利用すると参照だけが溜まってメモリが増えません。
rssDeltaMBはprocess.memoryUsage().rssの差をMBで書いたものです。- よくある間違いは、
readable.pipe(writable)で終わらせることです。背圧は守ってくれますが、 エラーを上に伝えないので、消費者が死んでも生産者は読み続けます。
このバージョンのデフォルト値を自分で測る
/root/work/backpressure/flow.mjsにmakeSink({hwm, delayMs})をexportし、/root/work/backpressure/report.jsonにnodeとhwm.writableDefault・hwm.objectMode・hwm.readableDefaultを測って書いてください。
new Writable({write(c, e, cb) { cb(); }}).writableHighWaterMarkを出力してみてください。覚えている数字と違うかもしれません。
makeSinkは、受け取ったバイト数と断片の数をreceivedBytes・receivedChunksに数えておき、delayMsが0より大きければ、その分だけ遅れてコールバックを呼びます。
write()がfalseを返す地点
fillUntilFalse(stream, chunk)をexportしてください。falseが出るまでwrite()を呼び、falseを返したその呼び出しまで含めた回数を返します。
write()は「この断片を受け取って入れたあと、溜まった量が上限より少ないか」を返します。上限と同じになった瞬間からfalseです。
そのため、64KBの上限に64KBを1回書くと、最初の呼び出しがすでにfalseです。この数字が実感できると、次のステップの待機ルールが自然になります。
止まって、drainを待つ
writeAll(stream, chunks)をexportしてください。順序を守ってすべて書きますが、write()がfalseならdrainを待ってから続きを書き、{peakBuffered}を返します。
await new Promise(r => stream.once("drain", r))の1行が背圧のすべてです。
onではなくonceを使ってください。毎回onで付けるとリスナーが溜まって警告が出て、結局それもメモリです。
戻り値を無視すると、どこまで溜まるのか
floodNoWait(stream, chunks)をexportし、64KBの断片2000個を戻り値を無視したまま流し込んで、runs.floodにchunks・chunkBytes・hwm・peakBuffered・bufferedMB・rssDeltaMB・wallMs・receivedChunksを書いてください。
peakBufferedはstream.writableLengthの最大値です。流し込んだ分がそのまま溜まります。上限は何も止めません。
断片を配列であらかじめ作ると、測る前にメモリをすでに使い切った状態になります。ジェネレーターで送るたびに作ってください。rssDeltaMBがほぼ0なら、その落とし穴にはまっています。
同じ量を背圧を守って送る
同じ断片数とサイズをwriteAllで送ってruns.pacedに同じ項目を書き、backpressureにbufferRatio(flood/pacedのpeakBuffered、小数第1位)、rssRatio(小数第1位)、timeRatio(paced/floodのwallMs、小数第2位)を書いてください。
溜まった量は数千倍の差が出るのに、timeRatioは1付近です。
消費者が受け取る速度は、どちらでも同じだからです。待つことが遅いわけではありません。 待たなければ、その差がメモリに溜まるだけです。
失敗したときに誰が片付けるのか
pipeThrough(readable, writable)をexportしてください。stream/promisesのpipelineを使い、すべて流し終えれば完了し、どちらかが失敗したらそのエラーで拒否される必要があります。同じ量を流してruns.pipelineにchunks・chunkBytes・hwm・peakBuffered・wallMs・receivedChunksを書いてください。
pipe()も背圧は守ってくれます。ところが消費者が死んでも何も起きないため、生産者は読み続け、エラーは誰も聞いていない場所で爆発します。
pipelineは、片方が失敗すると残りを片付け、エラーを呼び出した側まで伝えます。採点ツールは、わざと失敗する消費者をつないでそれを確かめます。
何秒後に壊れるのか
estimateQueueBytes({producerBps, consumerBps, seconds})とsecondsUntil({producerBps, consumerBps, limitBytes})をexportしてください。消費が生産より速いか同じなら、それぞれ0とInfinityです。数値でない値が入ってきたら例外を投げます。
溜まる速度はmax(0, 생산 - 소비)で(プレースホルダーは生産速度と消費速度です)、上限までの残り時間は한도 / 그 속도です(プレースホルダーは上限と溜まる速度です)。算数がすべてです。
この2行があれば、「メモリが少し上がっているのですが」ではなく「いまの速度だと4分後に上限に達します」と言えます。容量計画はここから始まります。