TT Lab
开始
学习 学习路径 课程

实时通信 — WebSocket、gRPC 流式调用与 WebRTC

gRPC 流式调用的四种形态

在 TT Lab 中继续学习

一句话总结

一次 gRPC 调用就是一个 HTTP/2 流,消息以前面带有长度的片段在其上流动,而状态码则放在最末尾的 trailer 中送来。

为什么需要它

如果只用一元调用来构建实时功能,就变成了轮询。像语音识别的部分结果、LLM 的令牌、行情这种一点点生成的结果,如果攒在一起一次发送,那么收到第一个片段所需的时间就等于整个处理时间。gRPC 核心概念文档 把调用分成四种形态——一元、服务器流式、客户端流式、双向流式。形态取决于 proto 的 rpc 声明中,请求和响应哪一侧加了 stream,生成的代码会创建与该形态相匹配的存根(stub)。

工作原理

查看 gRPC over HTTP/2 规范 就会发现,一次调用就是一个 HTTP/2 流。请求以 :path /패키지.서비스/메서드(占位符依次为包名、服务名、方法名)、content-type: application/grpc 开始,如果有截止时间,还有 grpc-timeout 头,消息则按 1 字节压缩标志 + 4 字节长度 + protobuf 字节依次连接。响应也是在相同形态的消息之后,以 trailer 中的 grpc-status 和 grpc-message 结束。

这种结构带来三点结论。

第一,只有边生成边发送,流式才是真正的流式。在 Python 服务器中用生成器 yield,gRPC 会逐个取出并发送。如果把结果全部攒进 list 再返回,第一条消息就会与最后一条消息在同一时刻到达。代码的形态是流式,用户体验却是一元。

第二,状态在最后到来。如果服务器发送了三条消息之后以错误结束,客户端要在收到这三条消息之后才会看到错误。在 Python 客户端中,会在迭代过程中抛出 grpc.RpcError。如果用 list(stub.Method(...)) 这一行来接收,就会丢掉已经收到的那三条。如果是部分结果有意义的流,就必须在循环内逐个保存。状态码有 17 种,在实时服务中常见的是 DEADLINE_EXCEEDED(4)、CANCELLED(1)、UNAVAILABLE(14)、RESOURCE_EXHAUSTED(8)、ABORTED(10)、INVALID_ARGUMENT(3)。

第三,双向流的两个方向相互独立。服务器可以在收完请求之前就作答,并且正如核心概念文档所说,两个流可以按任意顺序读取和写入。正因为有这种自由,很容易出现互相等待而停住的死锁。客户端收到答复之后才发送下一个请求,而服务器如果用 list(request_iterator) 先把请求全部攒起来,那么客户端在等答复,服务器在等请求的结束(half-close)。如果没有截止时间,就会永远等下去。

截止时间要设置在调用上。截止时间指南 说明默认值实际上是无限,所以建议为每个调用都显式设置截止时间。截止时间通过 grpc-timeout 头传给服务器,一旦超过截止时间,客户端会收到 DEADLINE_EXCEEDED。但是服务器的处理线程并不会自动停止。服务器代码必须通过 context.is_active() 或 context.time_remaining() 自行确认。如果不确认,就会为了一个没有人会接收的结果,一直占用 CPU 和 DB 连接。

消息大小也有上限。大多数实现把接收消息的默认上限设为 4MB,超过时以 RESOURCE_EXHAUSTED 结束调用。与其把一个大结果装进一条消息,不如把它切成流来发送,这是流式的另一个用处。切开发送后,接收方可以从第一个片段开始处理,流量控制一次扣住的内存也缩小到片段的大小。反过来,如果片段切得太碎,每条消息都要附带的 5 字节头和序列化成本就会变大,所以像语音这样有节拍的数据,要按 20ms 这样自然的单位来切分。

在现场相遇的样子

“改成了流式,却仍然是一次性到达”,是服务器把结果攒在一起再发送,或者中间的代理、负载均衡器缓冲了响应。要区分这两种情况,只需在服务器旁边测一次、在用户一侧再测一次第一条消息的到达时刻。“客户端已经超时,服务器 CPU 却一直在转”,说明服务器代码没有检查截止时间。在超过截止时间的请求集中出现的瞬间,服务器已经被被抛弃的工作塞满,并蔓延成连新请求也变慢的连锁故障。

下一项实验要做什么

从 meter.proto 生成代码,并实现服务器、客户端和双向流式。评分器通过第一条消息的到达时刻判断是否是流式,通过逐行收发、截止时间为 3 秒的调用判断是否死锁。还要构建在流中途出错时能保住已收到的值的客户端,以及超过截止时间会自行停止的服务器。