实时通信 — WebSocket、gRPC 流式调用与 WebRTC
亲手验证 gRPC 流式调用的四种形态与截止时间
目标
从 meter.proto 生成代码并实现服务器、客户端和双向流式,然后确认流中途的错误和超过截止时间,在客户端和服务器上各自是什么样子。
为什么重要
gRPC 流式调用是在一个 HTTP/2 流上承载多条消息的方式。因此响应状态在最后到来,截止时间通过头部传到服务器,而双方如果互相等待就会停住。如果不了解这三点,就会出现“改成了流式,却仍然是一次性到达”“一出错,已经收到的结果就全部消失”“已经超时了,服务器 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,就会生成 meter_pb2.py(消息)和 meter_pb2_grpc.py(服务骨架和存根)。打开这两个文件,确认 Countdown、Sum、Chat 分别被生成为哪种形态的调用。
- 服务器流式——边生成边发送——在 /root/rt/grpc/server.py 中创建继承 meter_pb2_grpc.MeterServicer 的 Meter,以及 serve(port)。serve 把 Meter 挂到有 16 个线程的 grpc.server 上,在 127.0.0.1:port 监听,并一直等到结束。Countdown 从 request.n 到 1 逐个发送 Tick,每发送一个就休息 request.interval_ms 毫秒。评分器会测量第一个 Tick 什么时候到达。
- 客户端流式——读到底之后只应答一次——给 Meter 添加 Sum。把传入的 Num 流读到底,把个数和总和以 Total(count, sum) 返回一次。如果流什么都没有来就结束了,则是 Total(count=0, sum=0)。
- 双向流式——在收完之前就作答——给 Meter 添加 Chat。对每个传入的 Line,立即返回把 text 转为大写的 Line。评分器会发送一行、收到答复之后才发送下一行,并设置 3 秒的截止时间。
- 流中途的错误不会抹掉已经收到的内容——修改 Countdown,当 request.fail_at 不为 0 时,发送 fail_at 个 Tick 之后,用 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)。以 timeout 秒的截止时间调用 Countdown,返回由收到的 value 的 list、状态码名称,以及调用所花秒数构成的 tuple。截止时间必须通过存根调用的 timeout 参数来设置。
- 超过截止时间之后,服务器的工作也要停止——给 Meter 添加 Work 和 GetStats。Work 要工作 request.steps 次、每次 request.step_ms 毫秒,但每一步中如果 context.is_active() 为假就停止,并以 WorkDone(steps_done) 返回完成的步数。服务器要统计到目前为止所有 Work 完成步数的累计值,GetStats 以 Stats(work_steps) 返回该累计值。评分器会以 0.35 秒的截止时间调用一个每步 100ms、共 20 步的工作,然后等待 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 自己启动。
- 有两个常见错误:在双向流式中先用 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,就会生成 meter_pb2.py(消息)和 meter_pb2_grpc.py(服务骨架和存根)。打开这两个文件,确认 Countdown、Sum、Chat 分别被生成为哪种形态的调用。
rpc 声明中的 stream 是加在请求一侧还是响应一侧,决定了四种形态的区别。生成代码中的 unary_stream、stream_unary、stream_stream 就是这些形态。不要修改 proto 文件。
服务器流式——边生成边发送
在 /root/rt/grpc/server.py 中创建继承 meter_pb2_grpc.MeterServicer 的 Meter,以及 serve(port)。serve 把 Meter 挂到有 16 个线程的 grpc.server 上,在 127.0.0.1:port 监听,并一直等到结束。Countdown 从 request.n 到 1 逐个发送 Tick,每发送一个就休息 request.interval_ms 毫秒。评分器会测量第一个 Tick 什么时候到达。
在生成器中 yield,gRPC 会逐个发出。如果全部攒进 list 再 return,第一条消息就会与最后一条消息在同一时刻到达,这样一来,做成流式就没有意义了。
客户端流式——读到底之后只应答一次
给 Meter 添加 Sum。把传入的 Num 流读到底,把个数和总和以 Total(count, sum) 返回一次。如果流什么都没有来就结束了,则是 Total(count=0, sum=0)。
请求迭代器要等客户端关闭流(half-close)才会结束。等待结束,就是这种形态的契约。
双向流式——在收完之前就作答
给 Meter 添加 Chat。对每个传入的 Line,立即返回把 text 转为大写的 Line。评分器会发送一行、收到答复之后才发送下一行,并设置 3 秒的截止时间。
如果像 list(request_iterator) 这样先把请求全部攒起来,客户端就会等答复,服务器则等请求结束,互相停住。要在遍历迭代器的循环内直接 yield。
流中途的错误不会抹掉已经收到的内容
修改 Countdown,当 request.fail_at 不为 0 时,发送 fail_at 个 Tick 之后,用 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。评分器会针对基准服务器调用它。
状态码放在消息之后的 trailer 中。所以错误会在迭代过程中以 grpc.RpcError 抛出,而在此之前收到的消息已经在你手里了。如果用一次 list(...) 来接收,就会丢掉它们。e.code().name 就是状态码名称。
截止时间设置在调用上
在 /root/rt/grpc/client.py 中添加 countdown_deadline(target, n, interval_ms, timeout)。以 timeout 秒的截止时间调用 Countdown,返回由收到的 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) 返回该累计值。评分器会以 0.35 秒的截止时间调用一个每步 100ms、共 20 步的工作,然后等待 2.5 秒,查看累计值。
客户端收到截止时间超限,并不会让服务器的线程停止。服务器代码必须自行确认。如果不确认,就会为了一个没有人会接收的结果,一直占用 CPU 和 DB 连接。多个线程会累加累计值,所以要使用锁。