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

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

遅い購読者を切り離してライブ配信を守る

TT Labで続きを見る

目標

キューのポリシーとACKの期限、接続の寿命を実装して、実際のTCPで遅い購読者を分離します。

なぜ重要なのか

購読者1人の遅れが、全体の発行を止めたり、メモリを使い続けさせたりすることがあります。データの損失の許容範囲と、業務の確認を区別し、キャンセルのときにもリソースを回収するプログラムを作ります。Pythonの関数・例外・async/awaitの基礎が必要です。標準ライブラリだけを使い、インストールやインターネットは必要ありません。

ステップ

  1. 際限なく積み上がるキューに上限をかける。delivery.pyにmake_queue(capacity)を実装してください。空のasyncio.Queueを返し、maxsizeは入力と同じである必要があります。正のintだけを許可し、bool・実数・文字列・0・負の数はValueErrorです。
  2. 待つが、キャンセルされた発行は残さない。async offer_wait(queue, item)を追加してください。空きができるまで待って入れ、Noneを返します。待機中のキャンセルはCancelledErrorとして伝え、既存の項目と順序を保存します。キャンセルされた項目を、後で挿入しません。
  3. 古い座標を捨てて、現在の位置を残す。同期関数offer_latest(queue, item)を追加してください。空きがあれば挿入してNone、いっぱいなら、最も長く待った項目を1つ取り除いて新しい項目を入れた後、取り除いた項目を返します。捨てた項目のtask_doneも、1回対応させます。
  4. 止まった購読者を、ポリシーの例外で知らせる。ExceptionのサブクラスのSlowConsumerと、同期のoffer_disconnect(queue, item)を実装してください。空きがあれば挿入してNoneを返します。いっぱいなら、既存のキューを変えずにSlowConsumerを出します。接続を閉じることは、この関数の責任ではありません。
  5. 重複は飛ばし、欠落は隠さない。ExceptionのサブクラスのGapとCursor(last=-1)を追加してください。lastはインスタンスごとの公開属性で、-1–2147483647の範囲のintです。accept(seq)は、0–2147483647のintだけを受け取ります。seq<=lastならFalseで状態を維持、seq==last+1ならlastを更新してTrue、それより大きい連番ならGapを出して状態を維持します。boolなどの誤った型・範囲はValueErrorです。メモリのカーソルであり、永続的な業務処理は実装しません。
  6. 送信と業務の確認に、時間を1回だけ与える。async send_one(reader, writer, seq, timeout, clock=None)を作ってください。seqはステップ5と同じ有効な連番で、timeoutはboolを除く正の有限なint/floatで、誤っていればValueErrorです。デフォルトの時計はtime.monotonic、注入する時計は引数なしで秒を返します。ASCIIのEVENT連番の後に改行を、writer.writeで1回書き、drainを待った後、reader.readlineで、まったく同じ連番のACKと改行を受け取って初めてNoneです。EOF・誤ったACKはConnectionError、全体の期限超過はTimeoutErrorです。drainとACKは同じ締切時刻を共有し、完了後にも期限を確認します。この関数はwriterを閉じず、エラー・キャンセルを伝えます。
  7. 空でも失敗しても、接続を回収する。async serve_queue(reader, writer, queue, timeout)を追加してください。正の有限なtimeoutが渡されます。繰り返しgetしたseqをsend_oneで処理し、その呼び出しが終わるか失敗・キャンセルされた後にだけ、引き受けた項目にtask_doneを、ちょうど1回呼びます。空のキューの待機を含め、すべての終了経路で、writer.closeの後、wait_closedをtimeoutの中で待ちます。終了待ちのタイムアウト・ConnectionErrorには、writer.transport.abortで回収し、元の送信エラー・キャンセルを隠してはいけません。まだキューに残った項目の後始末は、呼び出し側の責任です。
  8. 遅い1人を切り離して、ライブ配信を続ける。同期のbroadcast(queues, item, disconnect)を完成させてください。queuesは名前→キューのディクショナリで、巡回のスナップショットの各キューに、offer_disconnectを適用します。SlowConsumerだった対象だけをディクショナリから取り除き、disconnect(name)を1回呼んだ後、他の購読者に引き続き配信します。他の例外は伝え、正常な返却はNoneです。コールバックは、例外なしで、該当のワーカーのキャンセル・接続の回収・残りのキューの後始末を担当します。実際のTCP検査では、slowは最初のACKを保留し、fastは毎回確認します。0–8のイベントを、fastがすべて受け取り、slowだけが切り離される必要があります。

参考

すべての関数は、/root/realtime/delivery.py 1つに置きます。mkdir -p /root/realtimeで作業フォルダーを作ってください。例の、まだ実装していない関数は枠のままにして、完了した関数を上書きしないでください。採点は前のステップの契約も検査します。1ステップの実行は5秒の制限で、学習者が書く時間の制限ではありません。TCPサーバー・一時ポート・購読者は、検査器が準備して終了します。このラボは、カーネルバッファの飽和・インターネットの性能・永続的な配信・実際の再接続の復旧を保証しません。ラボのセッションが終わるとファイルは保持されないので、必要なコードは、終了前に別に保管してください。

際限なく積み上がるキューに上限をかける

delivery.pyにmake_queue(capacity)を実装してください。空のasyncio.Queueを返し、maxsizeは入力と同じである必要があります。正のintだけを許可し、bool・実数・文字列・0・負の数はValueErrorです。

asyncio.Queueでは、0は無制限です。型の検査と範囲の検査を分けて、標準ライブラリを使ってください。

待つが、キャンセルされた発行は残さない

async offer_wait(queue, item)を追加してください。空きができるまで待って入れ、Noneを返します。待機中のキャンセルはCancelledErrorとして伝え、既存の項目と順序を保存します。キャンセルされた項目を、後で挿入しません。

put_nowaitとawait putの飽和時の動作は違います。キャンセルを、正常な成功に変えないでください。

古い座標を捨てて、現在の位置を残す

同期関数offer_latest(queue, item)を追加してください。空きがあれば挿入してNone、いっぱいなら、最も長く待った項目を1つ取り除いて新しい項目を入れた後、取り除いた項目を返します。捨てた項目のtask_doneも、1回対応させます。

容量2に10,11があるときに12が来たら、11,12が残ります。処理中の項目は、このキューから取り除けません。

止まった購読者をポリシーの例外で知らせる

ExceptionのサブクラスのSlowConsumerと、同期のoffer_disconnect(queue, item)を実装してください。空きがあれば挿入してNoneを返します。いっぱいなら、既存のキューを変えずにSlowConsumerを出します。接続を閉じることは、この関数の責任ではありません。

1つの関数は飽和のポリシーを判断し、接続の所有者が回収します。黙って捨てたり待ったりするのは、別のポリシーです。

重複は飛ばし、欠落は隠さない

ExceptionのサブクラスのGapとCursor(last=-1)を追加してください。lastはインスタンスごとの公開属性で、-1–2147483647の範囲のintです。accept(seq)は、0–2147483647のintだけを受け取ります。seq<=lastならFalseで状態を維持、seq==last+1ならlastを更新してTrue、それより大きい連番ならGapを出して状態を維持します。boolなどの誤った型・範囲はValueErrorです。メモリのカーソルであり、永続的な業務処理は実装しません。

重複の検査とギャップの検査の順序が重要です。7の次に9が来たら、8を処理した根拠がありません。

送信と業務の確認に、時間を1回だけ与える

async send_one(reader, writer, seq, timeout, clock=None)を作ってください。seqはステップ5と同じ有効な連番で、timeoutはboolを除く正の有限なint/floatで、誤っていればValueErrorです。デフォルトの時計はtime.monotonic、注入する時計は引数なしで秒を返します。ASCIIのEVENT連番の後に改行を、writer.writeで1回書き、drainを待った後、reader.readlineで、まったく同じ連番のACKと改行を受け取って初めてNoneです。EOF・誤ったACKはConnectionError、全体の期限超過はTimeoutErrorです。drainとACKは同じ締切時刻を共有し、完了後にも期限を確認します。この関数はwriterを閉じず、エラー・キャンセルを伝えます。

asyncio.wait_forには、毎回残りの時間を与えます。テストは時計を注入して、drain 0.8秒とACK 0.3秒が、全体の1秒を超えるかを確認します。

空でも失敗しても、接続を回収する

async serve_queue(reader, writer, queue, timeout)を追加してください。正の有限なtimeoutが渡されます。繰り返しgetしたseqをsend_oneで処理し、その呼び出しが終わるか失敗・キャンセルされた後にだけ、引き受けた項目にtask_doneを、ちょうど1回呼びます。空のキューの待機を含め、すべての終了経路で、writer.closeの後、wait_closedをtimeoutの中で待ちます。終了待ちのタイムアウト・ConnectionErrorには、writer.transport.abortで回収し、元の送信エラー・キャンセルを隠してはいけません。まだキューに残った項目の後始末は、呼び出し側の責任です。

項目の所有権のtry/finallyと、接続の所有権のtry/finallyは、範囲が違います。ACKの前にtask_doneを呼ぶと、joinが先に解けます。

遅い1人を切り離して、ライブ配信を続ける

同期のbroadcast(queues, item, disconnect)を完成させてください。queuesは名前→キューのディクショナリで、巡回のスナップショットの各キューに、offer_disconnectを適用します。SlowConsumerだった対象だけをディクショナリから取り除き、disconnect(name)を1回呼んだ後、他の購読者に引き続き配信します。他の例外は伝え、正常な返却はNoneです。コールバックは、例外なしで、該当のワーカーのキャンセル・接続の回収・残りのキューの後始末を担当します。実際のTCP検査では、slowは最初のACKを保留し、fastは毎回確認します。0–8のイベントを、fastがすべて受け取り、slowだけが切り離される必要があります。

最初の切り離しの後で、returnしないでください。検査器が実際のサーバーと2つのクライアントを起動するので、自分でサーバーを常時実行する必要はありません。