リアルタイム通信 — WebSocket・gRPC ストリーミング・WebRTC
gRPC ストリーミングの 4 つの形と期限を手で確かめる
目標
meter.protoからコードを生成し、サーバー・クライアント・双方向ストリーミングを実装したあと、ストリームの途中のエラーと期限超過が、クライアントとサーバーでそれぞれどう見えるかを確認します。
なぜ重要なのか
gRPCストリーミングは、HTTP/2ストリーム1つの上に、メッセージを複数載せる方式です。そのため、応答のステータスは最後に来て、期限はヘッダーでサーバーまで伝わり、両側がお互いを待つと止まります。この3つを知らないと、「ストリーミングに変えたのに一度に来る」「エラーが出ると受け取った結果が全部消える」「タイムアウトしたのにサーバーのCPUが回り続ける」のようなことが起こります。音声AIの途中の認識結果とトークンストリーミングが、まさにこの形で流れます。
ステップ
- 契約からコードを生成する: /opt/fixtures/rt/grpc/meter.protoからPythonコードを生成して、/root/rt/grpcに置いてください。/opt/rt-lab/bin/python -m grpc_tools.protoc -I/opt/fixtures/rt/grpc --python_out=/root/rt/grpc --grpc_python_out=/root/rt/grpc /opt/fixtures/rt/grpc/meter.protoの1行で、meter_pb2.py(メッセージ)とmeter_pb2_grpc.py(サービスの骨組みとスタブ)ができます。2つのファイルを開いて、Countdown・Sum・Chatがどの形の呼び出しとして作られたかを確認してください。
- サーバーストリーミング(できあがるそばから送る): /root/rt/grpc/server.pyに、meter_pb2_grpc.MeterServicerを継承したMeterとserve(port)を作ってください。serveは、スレッド16個のgrpc.serverにMeterを付けて、127.0.0.1:portで待ち受け、終わるまで待機します。Countdownは、request.nから1までTickを1つずつ送り、1つ送るたびにrequest.interval_msミリ秒休みます。採点ツールは、最初のTickがいつ届くかを測ります。
- クライアントストリーミング(最後まで受け取って1回答える): MeterにSumを追加してください。入ってくるNumストリームを最後まで読み、個数と合計をTotal(count, sum)として1回返します。何も来ずに終わったストリームなら、Total(count=0, sum=0)です。
- 双方向ストリーミング(全部受け取る前に答える): MeterにChatを追加してください。入ってくるLineごとに、textを大文字に変えたLineをすぐに返します。採点ツールは、1行送って答えを受け取ってからでないと、次の行を送りません。そして、3秒の期限をかけます。
- ストリームの途中のエラーは、すでに受け取ったものを消さない: Countdownを直して、request.fail_atが0でなければ、Tickをfail_at個送ったあと、context.abort(grpc.StatusCode.ABORTED, "fail_at")で終わらせてください。そして、/root/rt/grpc/client.pyにread_countdown(target, n, interval_ms, fail_at, timeout=10)を作ってください。targetにチャネルを開いてCountdownを呼び、受け取ったvalueのlistと、終了したステータスコード名(正常なら"OK")をtupleで返します。採点ツールは、基準サーバーを相手に呼びます。
- 期限は呼び出しにかける: /root/rt/grpc/client.pyにcountdown_deadline(target, n, interval_ms, timeout)を追加してください。Countdownを期限timeout秒で呼び、受け取ったvalueのlist、ステータスコード名、呼び出しにかかった秒数をtupleで返します。期限は、スタブ呼び出しのtimeout引数でかける必要があります。
- 期限が過ぎたらサーバーの仕事も止まる: MeterにWorkとGetStatsを追加してください。Workは、request.steps回、request.step_msミリ秒ずつ作業しますが、ステップごとにcontext.is_active()が偽なら止まり、終えたステップ数をWorkDone(steps_done)として返します。サーバーは、これまでのすべてのWorkが終えたステップ数の累計を数えておき、GetStatsがその累計をStats(work_steps)として返します。採点ツールは、100msの20ステップの仕事を0.35秒の期限で呼んでから、2.5秒待って累計を見ます。
参考
- 作業フォルダは/root/rt/grpcです。mkdir -p /root/rt/grpcで、先に作ってください。
- サーバーは、cd /root/rt/grpc && /opt/rt-lab/bin/python -c "import server; server.serve(50051)"で起動してみられます。採点ツールは、空いているポートを選んで、別に起動します。
- クライアントのステップの相手である基準サーバーは、/opt/rt-lab/bin/python /opt/fixtures/rt/grpc/refserver.py --port 50052で、自分で起動できます。
- よくあるミスが2つあります。双方向ストリーミングでリクエストをlistで先にすべて集めて、お互いに止まることと、期限を呼び出しではなく、クライアント側の時計だけで測って、サーバーが期限を知らないままにすることです。
- Pythonは、必ず/opt/rt-lab/bin/pythonで実行します。このラボのライブラリは、その仮想環境にだけ入っていて、普通のpython3で実行すると、ModuleNotFoundErrorが出ます。alias rpy=/opt/rt-lab/bin/pythonのように短くしておくと便利です。
- ラボのPodは、外へ出る接続が塞がれています。すべての通信は、同じPodの中の127.0.0.1で行われ、インストールやダウンロードは必要ありません。
- 採点ツールは、コードを別プロセスで読み込んで、実際に接続を張ってみます。例のファイルは関数の枠にすぎないので、そのままでは合格しません。前のステップで完成させた関数は、消さないでください。
- ラボのセッションが終わると、/rootのファイルは残りません。必要なコードは、終える前に別に保管してください。
契約からコードを生成する
/opt/fixtures/rt/grpc/meter.protoからPythonコードを生成して、/root/rt/grpcに置いてください。/opt/rt-lab/bin/python -m grpc_tools.protoc -I/opt/fixtures/rt/grpc --python_out=/root/rt/grpc --grpc_python_out=/root/rt/grpc /opt/fixtures/rt/grpc/meter.protoの1行で、meter_pb2.py(メッセージ)とmeter_pb2_grpc.py(サービスの骨組みとスタブ)ができます。2つのファイルを開いて、Countdown・Sum・Chatがどの形の呼び出しとして作られたかを確認してください。
rpc宣言のstreamがリクエスト側に付いているか、レスポンス側に付いているかで、4つの形に分かれます。生成コードのunary_stream・stream_unary・stream_streamが、その形です。protoファイルは直しません。
サーバーストリーミング(できあがるそばから送る)
/root/rt/grpc/server.pyに、meter_pb2_grpc.MeterServicerを継承したMeterとserve(port)を作ってください。serveは、スレッド16個のgrpc.serverにMeterを付けて、127.0.0.1:portで待ち受け、終わるまで待機します。Countdownは、request.nから1までTickを1つずつ送り、1つ送るたびにrequest.interval_msミリ秒休みます。採点ツールは、最初のTickがいつ届くかを測ります。
ジェネレーターでyieldすると、gRPCが1つずつ送り出します。すべてlistに集めてreturnすると、最初のメッセージが最後のメッセージと同じ時刻に届き、それではストリーミングにした意味がありません。
クライアントストリーミング(最後まで受け取って1回答える)
MeterにSumを追加してください。入ってくるNumストリームを最後まで読み、個数と合計をTotal(count, sum)として1回返します。何も来ずに終わったストリームなら、Total(count=0, sum=0)です。
リクエストのイテレーターは、クライアントがストリームを閉じて(half-close)はじめて終わります。終わりを待つことが、この形の契約です。
双方向ストリーミング(全部受け取る前に答える)
MeterにChatを追加してください。入ってくるLineごとに、textを大文字に変えたLineをすぐに返します。採点ツールは、1行送って答えを受け取ってからでないと、次の行を送りません。そして、3秒の期限をかけます。
list(request_iterator)のようにリクエストを先にすべて集めると、クライアントは答えを待ち、サーバーはリクエストの終わりを待って、お互いに止まります。イテレーターを回すループの中で、すぐにyieldしてください。
ストリームの途中のエラーは、すでに受け取ったものを消さない
Countdownを直して、request.fail_atが0でなければ、Tickをfail_at個送ったあと、context.abort(grpc.StatusCode.ABORTED, "fail_at")で終わらせてください。そして、/root/rt/grpc/client.pyにread_countdown(target, n, interval_ms, fail_at, timeout=10)を作ってください。targetにチャネルを開いてCountdownを呼び、受け取ったvalueのlistと、終了したステータスコード名(正常なら"OK")をtupleで返します。採点ツールは、基準サーバーを相手に呼びます。
ステータスコードは、メッセージのあとのトレーラーに載って来ます。そのため、エラーは反復の途中でgrpc.RpcErrorとして飛び出し、その前に受け取ったメッセージは、すでに自分の手元にあります。list(...)1回で受け取ると、それを失います。e.code().nameがステータスコード名です。
期限は呼び出しにかける
/root/rt/grpc/client.pyにcountdown_deadline(target, n, interval_ms, timeout)を追加してください。Countdownを期限timeout秒で呼び、受け取ったvalueのlist、ステータスコード名、呼び出しにかかった秒数をtupleで返します。期限は、スタブ呼び出しのtimeout引数でかける必要があります。
timeoutを呼び出しにかけると、grpc-timeoutヘッダーでサーバーにも伝わります。クライアントで別に時計を測って止めると、サーバーは期限を知らないまま働き続けます。期限が過ぎたときのステータスはDEADLINE_EXCEEDEDです。
期限が過ぎたらサーバーの仕事も止まる
MeterにWorkとGetStatsを追加してください。Workは、request.steps回、request.step_msミリ秒ずつ作業しますが、ステップごとにcontext.is_active()が偽なら止まり、終えたステップ数をWorkDone(steps_done)として返します。サーバーは、これまでのすべてのWorkが終えたステップ数の累計を数えておき、GetStatsがその累計をStats(work_steps)として返します。採点ツールは、100msの20ステップの仕事を0.35秒の期限で呼んでから、2.5秒待って累計を見ます。
クライアントが期限超過を受け取っても、サーバーのスレッドは止まりません。サーバーコードが自分で確認する必要があります。確認しなければ、誰も受け取らない結果のために、CPUとDB接続を最後まで使います。複数のスレッドが累計を増やすので、ロックを使ってください。