TT Lab
はじめる
学ぶ 学習パス コース

遅い購読者がライブ配信を止めた

ライブ配信が遅れる理由

TT Labで続きを見る

一言でいうと

キューは処理すべき仕事を一時的に保管しますが、遅いコンシューマーを速くはしません。

なぜ必要なのか

猫の救助隊の位置をライブ配信すると想像してみましょう。購読者9人は新しい座標をすぐ画面に描きますが、1人は画面の処理を止めました。サーバーがすべての購読者の完了を順番に待つなら、止まった1人が、残り9人の地図まで止めてしまいます。逆に、いったん全員のキューに入れて次の仕事をすれば、パブリッシャーは楽になりますが、遅い人のキューは大きくなり続けます。待ちをなくしたのではなく、メモリに移しただけです。

このコースは、Pythonの関数・クラス・例外と、async/awaitの基礎を知っている学習者を対象にしています。前の「TCPの荷物」のコースを終えているなら、バイトの受信とメッセージの完了が違うことを思い出してください。今回は、メッセージが完成した後でも、業務が滞ることがあるという問題を扱います。実際の外部サービスを攻撃したり、インターネットを開いたりする必要はなく、同じラボのコンテナの中の、2つのTCP購読者で再現します。

どう動くのか

毎秒100個が入ってきて60個を処理するなら、捨てない待機キューは、処理率が変わらない間、毎秒約40個ずつ増えます。キューを4,000マスに増やしても、飽和する時点が遅れるだけで、継続的な不足は解決しません。キューの長さだけでなく、最も長く待っている項目の年齢も見る必要があります。座標100個が、1秒分なのか10分分なのかで、ユーザーが見る画面の意味が違います。

asyncio.Queue(maxsize=4)は、キューの中で待つ項目を4つに制限します。ワーカーがgetで取った項目は、その長さからは外れますが、ACKを待つ間もメモリに残っています。接続あたり待機4個と処理中1個、それに、シリアライズのバッファとソケットのバッファが別にあります。したがって、キューの長さ4を、全体のメモリ4項目の保証とは呼びません。項目のサイズも制限して初めて、バイト単位の予算を立てられます。

次の表で、パブリッシャーはプロデューサー、画面を更新するワーカーはコンシューマーです。

呼び出し 意味 よくある誤解
await queue.put(item) 空きができるまで待って入れる すべての購読者の処理が終わった
queue.put_nowait(item) すぐ入れるか、QueueFullを出す いっぱいでも黙って待つ
await queue.get() 待機中の項目1つの所有権を取る その項目の業務が成功した
queue.task_done() 取った作業の後始末を記録する リモートのデータが永続保存された

イベントループは、awaitで制御を手放している間に、他の作業を進めます。awaitなしで無限に繰り返す関数は、他の接続も止めます。async defという宣言そのものが、関数を並列のCPU作業に変えるわけではありません。またasyncio.Queueのmaxsize=0は、ゼロマスではなく無制限です。最初の段階で、正の整数だけを受け取るようにする理由です。boolはPythonでintのサブタイプですが、このAPIの容量としては拒否します。

いっぱいのキューで待っているputをキャンセルしたら、その発行も終了する必要があります。キャンセルを捕まえて成功のように返すと、呼び出し側は、入っていないイベントを送ったと信じてしまうことがあります。CancelledErrorを不用意に飲み込まず、既存の項目がそのまま残るかと、キャンセルされた項目が後で現れないかを、一緒に確認します。

現場での姿

ライブ字幕は、生成器が一時的に速くなることがあり、小さなキューが瞬間的な差を吸収します。しかし、モバイルの画面がずっと遅いなら、ある時点で、待機・省略・接続の終了のどれかを選ぶ必要があります。注文処理は、古い注文を捨てられないので、画面の座標と同じポリシーは使いません。キューの実装を選ぶ前に、データの意味を問うことが先です。

この授業のACKの保留は、購読者の業務処理が遅い状況です。カーネルの送信バッファの飽和や、インターネットのパケット損失を、直接測った実験ではありません。測った層をはっきりさせないと、キューの長さを変えて、ネットワーク障害まで解決したと勘違いしてしまいます。

次の確認ですること

空の容量2のキューにAとBを入れて、ワーカーがAを取ったと書いてみてください。qsizeは1ですが、task_doneをまだ呼んでいないので、未完了の作業は2つです。Cを入れるとqsizeは再び2になり、Dを入れようとするプロデューサーは待ちます。ワーカーがBを取る瞬間に空きができて、Dが入れます。Aの業務完了と、次の項目をキューに入れられる時点は、同じ出来事ではありません。

ここでAがいつまでもACKされなければ、ワーカーはBを取れません。キューにはBとCが残り、プロデューサーはDで止まります。この接続のキューだけを見れば、正常なバックプレッシャーです。しかし、そのプロデューサーが他のすべての購読者を担当しているなら、ライブ配信全体の停止に広がります。同じ道具でも、所有の範囲と待つ場所によって、障害の影響範囲が変わります。関数1つだけを見ず、誰がその関数を待っているかを、矢印で描いてみてください。

キューの上限、処理中の項目、キャンセルされたputの意味を、クイズで区別します。上の例で、AのACKが戻ってくる場合も追跡して、止まったプロデューサーがどんな順序で再び進むかを比べてください。最後のモジュールの統合ラボで、make_queueとoffer_waitを自分で実装した後、同じキューに3つのポリシーを適用します。