漏れているのではなく、行列ができているのだ
一言でいうと
速い生産者と遅い消費者をそのままつなぐと、残った差は必ずどこかに溜まり、Nodeではそのどこかがこのプロセスのヒープです。write()の戻り値はそれを防ぐためのつまみですが、戻り値を読んで止まるのは書く側の役目です。
なぜ必要なのか
「メモリリークがあるようです」で始まる調査の多くは、実はリークではありません。毎秒10MBを作り出す側と、毎秒2MBを受け取る側をつなぐと、余る8MBはどこかにあるはずで、そのどこかがストリームの内部バッファーです。ヒープスナップショットを取るとBufferが大量に確保されているのに、コードにはそれを握っている場所がありません。誰も漏らしておらず、ただ行列が長くなっただけです。
この事故が特に遅く発見されるのは、小さな入力では絶対に起きないからです。テストデータ2MBでは、バッファーが満ちる前に終わります。本番で800MBのファイルが1つ入ってくる日に初めて表に出て、その日プロセスはメモリの上限にぶつかって死にます。KubernetesではOOMKilledと記録され、再起動され、同じリクエストがまた入ってきます。
そのため、この事故はコードレビューで捕まえるべきです。捕まえる場所は1か所です。書く側がwrite()の戻り値を見ているか。 見ていないなら、そのコードがいつ壊れるかは入力サイズだけにかかっており、そのサイズは私たちが決めるものではありません。
どう動くのか
NodeのWritableストリームは内部にキューを持ち、highWaterMarkというしきい値を持ちます。write()は断片をキューに入れたあと、溜まった量がしきい値より少なければtrue、そうでなければfalseを返します。公式ドキュメントの表現をそのまま移すと、falseは「この断片は受け取ったが、これ以上書く前にdrainイベントを待ってほしい」という意味です(ストリームドキュメント)。
大切なのは、falseが書き込みを拒否しないという点です。無視して書き続ければ、ストリームは受け取り続けてキューに入れます。上限は何も止めません。止まるのは完全に書く側の責任で、それがこの2行です。
for (const chunk of chunks) {
if (!stream.write(chunk)) {
await new Promise((resolve) => stream.once("drain", resolve));
}
}
しきい値のデフォルト値はバージョンによって変わったことがあるので、暗記せずに測るほうがよいです。ラボイメージのNode 22.11.0で測ると、WritableもReadableも65536バイト(64KiB)で、オブジェクトモードは16個です。昔の記憶にある16384をそのまま使って計算すると食い違います。
違いは、測ると劇的です。同じイメージで64KBの断片2000個(125MB)を戻り値を無視したまま流し込むと、バッファーに125MBがそのまま溜まり、RSSが131MB増えました。同じ量をdrainを待ちながら送ると、溜まった最大は65536バイト、つまりしきい値そのままでした。数千倍の差です。
ところで、全体の時間はどうだったのでしょうか。背圧を守った側のほうが、むしろ少し速くなりました。消費者が受け取る速度は、どちらでも同じだからです。「待てば遅くなる」という直感は、ここでは成り立ちません。待たなければ、その差が時間ではなくメモリに行くだけです。
現場での姿
手で書いた背圧には、もう1つ空白があります。失敗したときに誰が片付けるのかです。readable.pipe(writable)は背圧は守ってくれますが、エラーを上に伝えてくれません。消費者側でエラーが起きても、生産者はそのことを知らないまま読み続け、エラーは誰も聞いていない場所で爆発します。ファイルハンドルとメモリが残るのはおまけです。
stream/promisesのpipelineは、その穴を埋めます。片方が失敗すると残りを片付け、エラーを呼び出した側まで伝えます。実際に測ってみると、失敗する消費者をつないだとき、pipelineはそのエラーで拒否され、生産者も一緒に片付けられましたが、pipeでつないだ側は何の知らせもなく永遠に待ち続けました。コードレビューでpipe(を探すことが、Syncを探すことと同じくらい価値がある理由です。
最後に、この話は算数に置き換えられるという点が、実務でいちばん役に立ちます。生産速度と消費速度がわかれば、毎秒溜まる量はその差で、上限までの残り時間は上限をその差で割った値です。「メモリが少し上がっているのですが」ではなく「いまの速度だと4分後に上限に達します」と言えれば、会話が推測から計画に変わります。
次のラボですること
このバージョンのデフォルトのしきい値を、自分で測って書くことから始めます。write()がfalseを返す地点を数え、drainを待つ書き込みを作り、同じ量を2つの方式で送って、溜まったバイトとRSSと全体の時間を並べて測ります。
そのあとpipelineで同じことをしながら、わざと失敗する消費者をつないで、エラーが上がってきて生産者が片付けられるかを確認します。最後に、溜まる速度と上限までの残り時間を計算する関数を作ります。採点ツールは、書き込まれた比率を2つの記録から再計算して照合し、作成した関数を実際に動かして同じ性質が出るかを確かめます。