ライブ配信が遅れる理由
一言でいうと
キューは処理すべき仕事を一時的に保管しますが、遅いコンシューマーを速くはしません。
なぜ必要なのか
猫の救助隊の位置をライブ配信すると想像してみましょう。購読者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つのポリシーを適用します。