实时通信 — WebSocket、gRPC 流式调用与 WebRTC
处理 gRPC 的截止时间、取消、流量控制、keepalive、重试与停机
目标
构建:把截止时间和取消传给后面服务的中继服务器、按慢消费者的节奏生产的流式服务器、声明了 keepalive 和重试策略的通道,以及收到 SIGTERM 时能优雅下线的服务器,并通过服务器一侧的数字来确认。
为什么重要
实时服务是由多个 gRPC 调用串成一条链的形态。链上有一处中断了截止时间或取消,就会把资源花在没有人在等待的工作上;绕开流量控制,一个慢消费者就能填满服务器内存;keepalive 对不上,安静的流就会被中间设备切断;对重试理解有误,就会把流的结果收两遍。每次部署都会断开流的服务器也很常见。每一种问题,用一行选项或一个回调就能防住,但如果不知道那一行,找原因就要花上好几天。
步骤
- 把收到的截止时间传给后面——首先用 /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。结果用 future.result() 等待。评分器会以 10 秒的截止时间发起一个每步 100ms、共 30 步的工作,在 0.4 秒时取消,2 秒之后查看后面的服务完成了多少步。
- 按慢消费者的节奏慢速生产——在 /root/rt/grpcc/server.py 中创建 Meter 和 serve(port)。Feed 发送 request.n 个 Chunk(seq, data),data 是 request.size 字节。GetStats 以 Stats(feed_produced) 返回所生成的 Chunk 数的累计值。评分器会请求 3000 个 64KiB 的块,只读取一个就停 3 秒,查看这期间服务器生成了多少个,以及服务器内存增长了多少。
- 在空闲超时和 ping 策略之间选择间隔——在 /root/rt/grpcc/client.py 中创建 make_channel(target)。返回设置了 keepalive 选项的 grpc.insecure_channel。评分器会设置一个会断开安静 3 秒的连接的中继器,在它后面放置一个会用 GOAWAY(too_many_pings) 拒绝频率高于 1 秒的 ping 的服务器,然后用你的通道发送一个在 7 秒内没有任何字节往来的调用。
- 用通道配置来声明重试——在 /root/rt/grpcc/client.py 中添加 make_retry_channel(target)。在 grpc.service_config 选项中,以 JSON 形式放入整个 rt.Meter 服务的 retryPolicy。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')" 这样启动。评分器会另行选择一个空闲端口来启动。
- 有两个常见错误:把没有截止时间的调用的 time_remaining() 原样传下去,以及把流中途断开的调用在应用程序中从头重新调用,从而把值收两遍。
- 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。结果用 future.result() 等待。评分器会以 10 秒的截止时间发起一个每步 100ms、共 30 步的工作,在 0.4 秒时取消,2 秒之后查看后面的服务完成了多少步。
截止时间通过头部传递,而取消不是。客户端取消时,只有前面服务器的 context 变为不活动,前面服务器阻塞着等待的后面调用并不知道。断开连接是前面服务器的职责。
按慢消费者的节奏慢速生产
在 /root/rt/grpcc/server.py 中创建 Meter 和 serve(port)。Feed 发送 request.n 个 Chunk(seq, data),data 是 request.size 字节。GetStats 以 Stats(feed_produced) 返回所生成的 Chunk 数的累计值。评分器会请求 3000 个 64KiB 的块,只读取一个就停 3 秒,查看这期间服务器生成了多少个,以及服务器内存增长了多少。
在生成器中逐个 yield 时,gRPC 只有在流量控制窗口允许时才取出下一个。如果另设生产线程,预先向无限制的队列中填充,就会绕开该窗口,使消费者的份额堆积在服务器内存中。
在空闲超时和 ping 策略之间选择间隔
在 /root/rt/grpcc/client.py 中创建 make_channel(target)。返回设置了 keepalive 选项的 grpc.insecure_channel。评分器会设置一个会断开安静 3 秒的连接的中继器,在它后面放置一个会用 GOAWAY(too_many_pings) 拒绝频率高于 1 秒的 ping 的服务器,然后用你的通道发送一个在 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 选项中,以 JSON 形式放入整个 rt.Meter 服务的 retryPolicy。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,也会原样上抛。如果把已经收到过值的流从头重新调用,就会把那些值收两遍。
部署期间也要把正在进行的调用做完
在 /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,不要在里面等待。