リアルタイム通信 — WebSocket・gRPC ストリーミング・WebRTC
gRPC の期限・キャンセル・フロー制御・keepalive・リトライ・終了を扱う
目標
期限とキャンセルを後段のサービスへ伝えるリレーサーバー、遅いコンシューマーに合わせて生産するストリーミングサーバー、keepalive・リトライポリシーを宣言したチャネル、SIGTERMでグレースフルに落ちるサーバーを作り、サーバー側の数字で確認します。
なぜ重要なのか
リアルタイムサービスは、複数のgRPC呼び出しがチェーンでつながった形です。チェーンの1か所が期限やキャンセルを切ると、誰も待っていない仕事にリソースが使われ、フロー制御を迂回すると、遅いコンシューマー1つがサーバーのメモリを埋め、keepaliveがずれると、静かなストリームが中間機器に切られ、リトライを誤解すると、ストリームの結果を2回受け取ります。デプロイのたびにストリームが切れるサーバーも、よくあります。それぞれは、オプション1行やコールバック1つで防げますが、その1行を知らなければ、原因を探すのに何日もかかります。
ステップ
- 受け取った期限を後ろへ渡す: まず、/opt/rt-lab/bin/python -m grpc_tools.protoc -I/opt/fixtures/rt/grpc --python_out=/root/rt/grpcc --grpc_python_out=/root/rt/grpcc /opt/fixtures/rt/grpc/meter.protoでコードを生成してください。そして、/root/rt/grpcc/front.pyにserve(port, backend)を作ってください。127.0.0.1:portでMeterサービスを待ち受け、Workだけを実装します。Workは、同じリクエストをbackendのアドレスのMeter.Workへ渡しながら、自分の呼び出しの残りの期限(context.time_remaining())を、そのままtimeoutとしてかけます。期限のない呼び出しなら、timeoutはかけません。後ろでエラーが出たら、同じステータスコードと説明でcontext.abortします。
- 前段がキャンセルされたら後段もキャンセルする: Workを直して、後段の呼び出しをstub.Work.future(...)で非同期にかけ、context.add_callbackで、自分の呼び出しが終わるかキャンセルされたときに、そのfutureをcancelするようにしてください。結果はfuture.result()で待ちます。採点ツールは、100msの30ステップの仕事を期限10秒でかけたあと、0.4秒でキャンセルし、2秒後に、後段のサービスが何ステップ進めたかを見ます。
- 遅いコンシューマーに合わせてゆっくり作る: /root/rt/grpcc/server.pyにMeterとserve(port)を作ってください。Feedは、request.n個のChunk(seq, data)を送り、dataはrequest.sizeバイトです。作ったChunk数の累計を、GetStatsがStats(feed_produced)として返します。採点ツールは、64KiBのものを3000個要求してから、1つだけ読んで3秒止まり、そのあいだにサーバーがいくつ作ったかと、サーバーのメモリがどれだけ増えたかを見ます。
- アイドルタイムアウトとpingポリシーの間で間隔を選ぶ: /root/rt/grpcc/client.pyにmake_channel(target)を作ってください。keepaliveオプションをかけたgrpc.insecure_channelを返します。採点ツールは、3秒のあいだ静かな接続を切る中継器を置き、その後ろに、1秒より頻繁なpingをGOAWAY(too_many_pings)で拒否するサーバーを置いたまま、7秒のあいだ何のバイトもやり取りされない呼び出しを、自分のチャネルで送ります。
- リトライはチャネル設定で宣言する: /root/rt/grpcc/client.pyにmake_retry_channel(target)を追加してください。grpc.service_configオプションに、rt.Meterサービス全体のretryPolicyをJSONで入れます。maxAttempts 4、initialBackoff "0.1s"、maxBackoff "1s"、backoffMultiplier 2、retryableStatusCodes ["UNAVAILABLE"]です。grpc.enable_retriesも1でオンにします。
- ストリームは最初の応答のあとは再試行しない: /root/rt/grpcc/client.pyにcollect(target, n, fail_at, fail_before)を追加してください。make_retry_channelでチャネルを開いて、Countdown(n=n, interval_ms=10, fail_at=fail_at, fail_before=fail_before, fail_unavailable=True)を呼び、受け取ったvalueのlistと、終了したステータスコード名を返します。アプリケーションから呼び直すことはしません。
- デプロイ中も実行中の呼び出しは終わらせる: /root/rt/grpcc/server.pyのserveにSIGTERMハンドラーを付けて、シグナルを受け取ったらserver.stop(3)を呼び、新しい呼び出しは拒否し、進行中の呼び出しには3秒を与えるようにしてください。Meterには、前のラボのように、nから1まで送るCountdownも置いてください。採点ツールは、300ms間隔の4個のCountdownを開始し、0.3秒でSIGTERMを送ったあと、そのストリームが最後まで届くか、新しい呼び出しが拒否されるか、プロセスが終了するかを見ます。
参考
- 作業フォルダは/root/rt/grpccです。mkdir -p /root/rt/grpccで、先に作ってください。
- 後段のサービスとクライアントのステップの相手は、基準サーバーです。/opt/rt-lab/bin/python /opt/fixtures/rt/grpc/refserver.py --port 50052 --min-ping-ms 1000で、自分で起動してみられます。
- 前段のサービスは、cd /root/rt/grpcc && /opt/rt-lab/bin/python -c "import front; front.serve(50061, '127.0.0.1:50052')"のように起動します。採点ツールは、空いているポートを選んで、別に起動します。
- よくあるミスが2つあります。期限のない呼び出しのtime_remaining()をそのまま渡すことと、ストリームの途中で切れた呼び出しを、アプリケーションで最初から呼び直して、値を2回受け取ることです。
- Pythonは、必ず/opt/rt-lab/bin/pythonで実行します。このラボのライブラリは、その仮想環境にだけ入っていて、普通のpython3で実行すると、ModuleNotFoundErrorが出ます。alias rpy=/opt/rt-lab/bin/pythonのように短くしておくと便利です。
- ラボのPodは、外へ出る接続が塞がれています。すべての通信は、同じPodの中の127.0.0.1で行われ、インストールやダウンロードは必要ありません。
- 採点ツールは、コードを別プロセスで読み込んで、実際に接続を張ってみます。例のファイルは関数の枠にすぎないので、そのままでは合格しません。前のステップで完成させた関数は、消さないでください。
- ラボのセッションが終わると、/rootのファイルは残りません。必要なコードは、終える前に別に保管してください。
受け取った期限を後ろへ渡す
まず、/opt/rt-lab/bin/python -m grpc_tools.protoc -I/opt/fixtures/rt/grpc --python_out=/root/rt/grpcc --grpc_python_out=/root/rt/grpcc /opt/fixtures/rt/grpc/meter.protoでコードを生成してください。そして、/root/rt/grpcc/front.pyにserve(port, backend)を作ってください。127.0.0.1:portでMeterサービスを待ち受け、Workだけを実装します。Workは、同じリクエストをbackendのアドレスのMeter.Workへ渡しながら、自分の呼び出しの残りの期限(context.time_remaining())を、そのままtimeoutとしてかけます。期限のない呼び出しなら、timeoutはかけません。後ろでエラーが出たら、同じステータスコードと説明でcontext.abortします。
前段で0.6秒の期限を受け取ったのに、後ろへは期限なしで呼び出すと、前段のクライアントが諦めたあとも、後段は最後まで働きます。残りの期限を渡すことが、期限の伝播です。期限のない呼び出しでtime_remaining()はとても大きな値になるので、そのまま渡さないでください。
前段がキャンセルされたら後段もキャンセルする
Workを直して、後段の呼び出しをstub.Work.future(...)で非同期にかけ、context.add_callbackで、自分の呼び出しが終わるかキャンセルされたときに、そのfutureをcancelするようにしてください。結果はfuture.result()で待ちます。採点ツールは、100msの30ステップの仕事を期限10秒でかけたあと、0.4秒でキャンセルし、2秒後に、後段のサービスが何ステップ進めたかを見ます。
期限はヘッダーで伝わりますが、キャンセルはそうではありません。クライアントがキャンセルすると、前段のサーバーのcontextだけが非アクティブになり、前段のサーバーが塞がって待っている後段の呼び出しは、知りません。接続を切ってあげるのは、前段のサーバーの役目です。
遅いコンシューマーに合わせてゆっくり作る
/root/rt/grpcc/server.pyにMeterとserve(port)を作ってください。Feedは、request.n個のChunk(seq, data)を送り、dataはrequest.sizeバイトです。作ったChunk数の累計を、GetStatsがStats(feed_produced)として返します。採点ツールは、64KiBのものを3000個要求してから、1つだけ読んで3秒止まり、そのあいだにサーバーがいくつ作ったかと、サーバーのメモリがどれだけ増えたかを見ます。
ジェネレーターで1つずつyieldすると、gRPCは、フロー制御ウィンドウが許すときにだけ次のものを取り出します。生産スレッドを別に置いて、制限のないキューにあらかじめ満たしておくと、そのウィンドウを迂回して、コンシューマーの分がサーバーのメモリにたまります。
アイドルタイムアウトとpingポリシーの間で間隔を選ぶ
/root/rt/grpcc/client.pyにmake_channel(target)を作ってください。keepaliveオプションをかけたgrpc.insecure_channelを返します。採点ツールは、3秒のあいだ静かな接続を切る中継器を置き、その後ろに、1秒より頻繁なpingをGOAWAY(too_many_pings)で拒否するサーバーを置いたまま、7秒のあいだ何のバイトもやり取りされない呼び出しを、自分のチャネルで送ります。
grpc.keepalive_time_msはping間隔、grpc.keepalive_timeout_msは答えを待つ時間、grpc.keepalive_permit_without_callsは呼び出しがないときにもpingするか、grpc.http2.max_pings_without_dataはデータなしで送るpingの上限(0なら無制限)です。間隔が長すぎれば中継器が、頻繁すぎればサーバーが切ります。
リトライはチャネル設定で宣言する
/root/rt/grpcc/client.pyにmake_retry_channel(target)を追加してください。grpc.service_configオプションに、rt.Meterサービス全体のretryPolicyをJSONで入れます。maxAttempts 4、initialBackoff "0.1s"、maxBackoff "1s"、backoffMultiplier 2、retryableStatusCodes ["UNAVAILABLE"]です。grpc.enable_retriesも1でオンにします。
リトライをアプリケーションのループで組むと、バックオフ・上限・どのステータスコードをリトライするかが、呼び出しごとにばらばらになります。INVALID_ARGUMENTのように、やり直しても同じ結果になるエラーは、リトライの一覧に入れません。
ストリームは最初の応答のあとは再試行しない
/root/rt/grpcc/client.pyにcollect(target, n, fail_at, fail_before)を追加してください。make_retry_channelでチャネルを開いて、Countdown(n=n, interval_ms=10, fail_at=fail_at, fail_before=fail_before, fail_unavailable=True)を呼び、受け取ったvalueのlistと、終了したステータスコード名を返します。アプリケーションから呼び直すことはしません。
リトライポリシーは、サーバーが応答ヘッダーや最初のメッセージを送る前に出たエラーだけを、再試行します。そのあとは、呼び出しが確定(committed)して、同じUNAVAILABLEでも、そのまま上がってきます。すでに受け取った値があるストリームを最初から呼び直すと、その値を2回受け取ります。
デプロイ中も実行中の呼び出しは終わらせる
/root/rt/grpcc/server.pyのserveにSIGTERMハンドラーを付けて、シグナルを受け取ったらserver.stop(3)を呼び、新しい呼び出しは拒否し、進行中の呼び出しには3秒を与えるようにしてください。Meterには、前のラボのように、nから1まで送るCountdownも置いてください。採点ツールは、300ms間隔の4個のCountdownを開始し、0.3秒でSIGTERMを送ったあと、そのストリームが最後まで届くか、新しい呼び出しが拒否されるか、プロセスが終了するかを見ます。
Kubernetesは、Podを落とすときにSIGTERMを送り、terminationGracePeriodSecondsのあとにSIGKILLを送ります。ハンドラーがなければ、PythonはSIGTERMですぐに死に、ストリーミング中だったクライアントは、すべてUNAVAILABLEを受け取ります。シグナルハンドラーの中では、stopを呼び出すだけにして、待たないでください。