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

リアルタイム通信 — WebSocket・gRPC ストリーミング・WebRTC

長く繋がる WebSocket ハブを運用する

TT Labで続きを見る

目標

websocketsで、発行・購読ハブと再接続クライアントを作り、遅い購読者・アイドルタイムアウト・黙って消えた相手・デプロイによる一斉切断に、順番に耐えられるようにします。

なぜ重要なのか

WebSocketサービスの障害は、ほとんどがプロトコルではなく、時間から来ます。1人が遅くなればハブのメモリがいっぱいになり、静かな接続はロードバランサーが先に切り、携帯電話がトンネルに入ると相手はcloseなしで消え、デプロイ1回で数万が同時に再接続します。このラボの5つのルール、つまり、購読者ごとの上限、1013での切断、アイドル上限より短いping、終了も一緒に待つこと、ジッター付きの再接続とseqでの続きの受け取りは、チャットでも相場でも、音声AIの途中結果でも、同じように必要です。

ステップ

  1. 再接続の間隔にジッターを混ぜる: /root/rt/wsops/client.pyにbackoff(attempt, base=0.2, cap=5.0, rng=random)を作ってください。attempt回目のリトライの前に待つ秒数を返します。値は、rng.uniform(0, min(cap, base * 2 ** attempt))を1回呼んで決めます(full jitter)。rngには、randomモジュールまたはrandom.Randomオブジェクトが来ます。
  2. クローズコードを見て再接続するか決める: /root/rt/wsops/client.pyにshould_reconnect(code)を追加してください。1001(去っていく)・1006(closeフレームなしで切断)・1011(サーバーの内部エラー)・1012(サービスの再起動)・1013(あとでもう一度)ならTrue、それ以外のコードはすべてFalseを返します。
  3. ハブが1人の発言を全員に運ぶ: /root/rt/wsops/hub.pyにrun(host, port, ring=1000, queue=100, ping_interval=1.0, ping_timeout=2.0)コルーチンを作ってください。websockets.asyncio.server.serveで、3つのパスを受け付けます。/pubに入ってきたテキストメッセージごとに、1から増えるseqを付けて、{"seq": n, "data": メッセージ}のJSON文字列にし、/subにつながったすべての購読者に送ります。/statsは、{"subscribers": 購読者数, "dropped": 切断した購読者数, "seq": 最後のseq}のJSONを1つ送って終わります。購読者ごとに、asyncio.Queueと、そのキューを空にしながら送るループを、別に持たせます。
  4. 遅い購読者1人だけを切る: 購読者キューのサイズをqueueで制限してください。キューがいっぱいで入れられなければ、その購読者をリストから外し、droppedを1上げて、クローズコード1013と理由"slow consumer"で閉じます。採点ツールは、読まない購読者1人とよく読む購読者1人をつないだまま、8000バイトのメッセージ8000個を発行し、よく読む側がすべて受け取るか、ハブのメモリがどれだけ増えたか、読まなかった側が1013を受け取るかを確認します。
  5. 静かな接続がロードバランサーに切られないようにする: serveに、ping_intervalとping_timeoutを、runの引数のまま渡してください。採点ツールは、2.5秒のあいだ静かな接続を切る中継器をハブの前に置き、pingを自分では送らない購読者で6秒待ってから、メッセージを1つ発行します。そのメッセージが購読者に届かなければなりません。
  6. 黙って消えた購読者を片づける: 購読者ループが、キューだけを待たずに、接続が閉じることも一緒に待つように直してください。ws.wait_closed()とq.get()をasyncio.wait(..., return_when=FIRST_COMPLETED)で一緒に待ち、接続が先に閉じたら、購読者リストから外してループを終えます。serveには、close_timeout=1.0も渡してください。採点ツールは、ハンドシェイクだけを行ってpingに答えない購読者をつないだあと、6秒以内に/statsのsubscribersが0に戻るかを確認します。
  7. 切れた位置から続きを受け取る: ハブに、最近のring個のメッセージを入れるリングバッファーを置き、/sub?last=Nでつながったら、seqがNより大きいメッセージを先に再送してから、リアルタイムで続けて送るようにしてください。N+1がバッファーからすでに押し出されていたら、{"reset": true, "seq": 最後のseq}を先に送ります。そして、/root/rt/wsops/client.pyにconsume(url, out_path, stop_after, base=0.1, cap=1.0)コルーチンを追加してください。urlにつないで受け取ったメッセージのseqを、1行に1つずつout_pathに追記し、接続が切れたら、should_reconnectとbackoffで待ってから、最後に書き込んだseqをlastとして渡して再接続します。seqがstop_afterに達したら終わります。採点ツールは、発行の途中で接続を2回切り、ファイルに1からstop_afterまでが、抜けなく1回ずつあるかを確認します。

参考

再接続の間隔にジッターを混ぜる

/root/rt/wsops/client.pyにbackoff(attempt, base=0.2, cap=5.0, rng=random)を作ってください。attempt回目のリトライの前に待つ秒数を返します。値は、rng.uniform(0, min(cap, base * 2 ** attempt))を1回呼んで決めます(full jitter)。rngには、randomモジュールまたはrandom.Randomオブジェクトが来ます。

ジッターがないと、一斉に切れたクライアント数万が、まったく同じ時刻に再び押し寄せて、生き返ったばかりのサーバーをまた倒します。上限(cap)がないと、10回目のリトライは数分後になります。採点ツールは、シードを決めたrandom.Randomを渡して、同じ値を計算してみます。

クローズコードを見て再接続するか決める

/root/rt/wsops/client.pyにshould_reconnect(code)を追加してください。1001(去っていく)・1006(closeフレームなしで切断)・1011(サーバーの内部エラー)・1012(サービスの再起動)・1013(あとでもう一度)ならTrue、それ以外のコードはすべてFalseを返します。

1008(ポリシー違反)や1002(プロトコルエラー)で切れた接続に再接続すると、同じ理由でまた切れます。リトライに意味があるのは、相手側の事情が変わりうる場合だけです。1000は、誰かが意図的に閉じたという意味です。

ハブが1人の発言を全員に運ぶ

/root/rt/wsops/hub.pyにrun(host, port, ring=1000, queue=100, ping_interval=1.0, ping_timeout=2.0)コルーチンを作ってください。websockets.asyncio.server.serveで、3つのパスを受け付けます。/pubに入ってきたテキストメッセージごとに、1から増えるseqを付けて、{"seq": n, "data": メッセージ}のJSON文字列にし、/subにつながったすべての購読者に送ります。/statsは、{"subscribers": 購読者数, "dropped": 切断した購読者数, "seq": 最後のseq}のJSONを1つ送って終わります。購読者ごとに、asyncio.Queueと、そのキューを空にしながら送るループを、別に持たせます。

発行する側が購読者ごとにawait sendを順番に呼ぶと、1人が遅いとき、その後ろの全員が待ちます。購読者ごとのキューは、その待ちを人ごとに切り離す仕掛けです。発行は、キューに入れるだけで、待ちません。

遅い購読者1人だけを切る

購読者キューのサイズをqueueで制限してください。キューがいっぱいで入れられなければ、その購読者をリストから外し、droppedを1上げて、クローズコード1013と理由"slow consumer"で閉じます。採点ツールは、読まない購読者1人とよく読む購読者1人をつないだまま、8000バイトのメッセージ8000個を発行し、よく読む側がすべて受け取るか、ハブのメモリがどれだけ増えたか、読まなかった側が1013を受け取るかを確認します。

キューに上限がないと、遅い購読者1人の分のメッセージが、ハブのメモリに際限なくたまります。メッセージを捨てることと接続を切ることのどちらがよいかは、データの性質が決めます。順序と欠落が重要なストリームなら、穴が空いたまま送り続けるより、切って再接続させるほうが正直です。

静かな接続がロードバランサーに切られないようにする

serveに、ping_intervalとping_timeoutを、runの引数のまま渡してください。採点ツールは、2.5秒のあいだ静かな接続を切る中継器をハブの前に置き、pingを自分では送らない購読者で6秒待ってから、メッセージを1つ発行します。そのメッセージが購読者に届かなければなりません。

websocketsの既定のping間隔は20秒なので、アイドル上限がそれより短い機器の後ろでは、静かな接続が先に切れます。よくあるAWS ALBの既定のアイドル上限は60秒で、社内プロキシはもっと短い場合が多いです。間隔は、最も短いアイドル上限より短くする必要があります。

黙って消えた購読者を片づける

購読者ループが、キューだけを待たずに、接続が閉じることも一緒に待つように直してください。ws.wait_closed()とq.get()をasyncio.wait(..., return_when=FIRST_COMPLETED)で一緒に待ち、接続が先に閉じたら、購読者リストから外してループを終えます。serveには、close_timeout=1.0も渡してください。採点ツールは、ハンドシェイクだけを行ってpingに答えない購読者をつないだあと、6秒以内に/statsのsubscribersが0に戻るかを確認します。

pingのタイムアウトでライブラリが接続を閉じても、キューを待つコルーチンは、新しいメッセージが来るまで起きません。そのあいだ、購読者リストには死んだ接続が残って、数字を膨らませ、メッセージを受け取ってためます。そして、答えのない相手にcloseを送ったあと、答えを待つ時間の既定値は10秒です。

切れた位置から続きを受け取る

ハブに、最近のring個のメッセージを入れるリングバッファーを置き、/sub?last=Nでつながったら、seqがNより大きいメッセージを先に再送してから、リアルタイムで続けて送るようにしてください。N+1がバッファーからすでに押し出されていたら、{"reset": true, "seq": 最後のseq}を先に送ります。そして、/root/rt/wsops/client.pyにconsume(url, out_path, stop_after, base=0.1, cap=1.0)コルーチンを追加してください。urlにつないで受け取ったメッセージのseqを、1行に1つずつout_pathに追記し、接続が切れたら、should_reconnectとbackoffで待ってから、最後に書き込んだseqをlastとして渡して再接続します。seqがstop_afterに達したら終わります。採点ツールは、発行の途中で接続を2回切り、ファイルに1からstop_afterまでが、抜けなく1回ずつあるかを確認します。

「受け取ったもの」ではなく、「処理して書き残したもの」を基準に再接続する必要があります。受け取っただけで書く前に切れると、そのメッセージは永遠に消えます。同じseqが2回来たら、1回だけ書いてください。