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

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

gRPC ストリーミングの 4 つの形と期限を手で確かめる

TT Labで続きを見る

目標

meter.protoからコードを生成し、サーバー・クライアント・双方向ストリーミングを実装したあと、ストリームの途中のエラーと期限超過が、クライアントとサーバーでそれぞれどう見えるかを確認します。

なぜ重要なのか

gRPCストリーミングは、HTTP/2ストリーム1つの上に、メッセージを複数載せる方式です。そのため、応答のステータスは最後に来て、期限はヘッダーでサーバーまで伝わり、両側がお互いを待つと止まります。この3つを知らないと、「ストリーミングに変えたのに一度に来る」「エラーが出ると受け取った結果が全部消える」「タイムアウトしたのにサーバーのCPUが回り続ける」のようなことが起こります。音声AIの途中の認識結果とトークンストリーミングが、まさにこの形で流れます。

ステップ

  1. 契約からコードを生成する: /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がどの形の呼び出しとして作られたかを確認してください。
  2. サーバーストリーミング(できあがるそばから送る): /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がいつ届くかを測ります。
  3. クライアントストリーミング(最後まで受け取って1回答える): MeterにSumを追加してください。入ってくるNumストリームを最後まで読み、個数と合計をTotal(count, sum)として1回返します。何も来ずに終わったストリームなら、Total(count=0, sum=0)です。
  4. 双方向ストリーミング(全部受け取る前に答える): MeterにChatを追加してください。入ってくるLineごとに、textを大文字に変えたLineをすぐに返します。採点ツールは、1行送って答えを受け取ってからでないと、次の行を送りません。そして、3秒の期限をかけます。
  5. ストリームの途中のエラーは、すでに受け取ったものを消さない: 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で返します。採点ツールは、基準サーバーを相手に呼びます。
  6. 期限は呼び出しにかける: /root/rt/grpc/client.pyにcountdown_deadline(target, n, interval_ms, timeout)を追加してください。Countdownを期限timeout秒で呼び、受け取ったvalueのlist、ステータスコード名、呼び出しにかかった秒数をtupleで返します。期限は、スタブ呼び出しのtimeout引数でかける必要があります。
  7. 期限が過ぎたらサーバーの仕事も止まる: 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秒待って累計を見ます。

参考

契約からコードを生成する

/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接続を最後まで使います。複数のスレッドが累計を増やすので、ロックを使ってください。