实时通信 — WebSocket、gRPC 流式调用与 WebRTC
运营长连接 WebSocket 中枢
目标
用 websockets 构建发布/订阅中心和重连客户端,使它依次承受慢订阅者、空闲超时、悄无声息消失的对端,以及因部署而导致的集中断开。
为什么重要
WebSocket 服务的故障大多不是出在协议上,而是出在时间上。一个人变慢,中心的内存就被填满;安静的连接会被负载均衡器先断开;手机进入隧道,对端就没有 close 地消失;一次部署,几万个连接就会同时重新连上。本实验的五条规则——按订阅者设置上限、用 1013 断开、使用比空闲上限更短的 ping、同时等待关闭、带抖动的重连以及用 seq 接续——无论是聊天、行情还是语音 AI 的部分结果,都同样需要。
步骤
- 给重连间隔掺入抖动——在 /root/rt/wsops/client.py 中创建 backoff(attempt, base=0.2, cap=5.0, rng=random)。返回第 attempt 次重试之前要等待的秒数。该值只通过一次 rng.uniform(0, min(cap, base * 2 ** attempt)) 来确定(full jitter)。rng 会传入 random 模块或 random.Random 对象。
- 根据关闭码决定是否重新连接——在 /root/rt/wsops/client.py 中添加 should_reconnect(code)。如果是 1001(离开)、1006(没有 close 帧就断开)、1011(服务器内部错误)、1012(服务重启)、1013(稍后重试),就返回 True,其他所有码都返回 False。
- 让中心把一个人说的话传给所有人——在 /root/rt/wsops/hub.py 中创建 run(host, port, ring=1000, queue=100, ping_interval=1.0, ping_timeout=2.0) 协程。用 websockets.asyncio.server.serve 接收三个路径。对于通过 /pub 进来的每条文本消息,加上从 1 开始递增的 seq,做成 {"seq": n, "data": 消息} 形式的 JSON 字符串,发送给通过 /sub 连上的所有订阅者。/stats 发送一个 {"subscribers": 订阅者数, "dropped": 被断开的订阅者数, "seq": 最后的 seq} JSON 之后结束。为每个订阅者单独设置一个 asyncio.Queue,以及一个清空该队列并发送的循环。
- 只断开一个慢订阅者——用 queue 限制订阅者队列的大小。如果队列已满、无法放入,就把该订阅者从列表中移除,dropped 加 1,然后用关闭码 1013 和原因 "slow consumer" 关闭。评分器会在接入一个不读取数据的订阅者和一个读取良好的订阅者的情况下,发布 8000 条 8000 字节的消息,查看读取良好的一方是否全部收到、中心的内存增长了多少,以及原来不读取的一方是否收到了 1013。
- 不让安静的连接被负载均衡器切断——把 ping_interval 和 ping_timeout 按 run 的参数原样传给 serve。评分器会在中心前面放置一个会断开安静 2.5 秒的连接的中继器,用一个不会自行发送 ping 的订阅者等待 6 秒之后,发布一条消息。该消息必须到达订阅者。
- 清除悄无声息消失的订阅者——修改订阅者循环,使它不只是等待队列,同时也等待连接关闭。用 asyncio.wait(..., return_when=FIRST_COMPLETED) 同时等待 ws.wait_closed() 和 q.get(),如果连接先关闭,就把它从订阅者列表中移除并结束循环。也要把 close_timeout=1.0 传给 serve。评分器会接入一个只做握手、不应答 ping 的订阅者,然后查看 6 秒内 /stats 的 subscribers 是否回到 0。
- 从断开的位置接续——在中心中设置一个存放最近 ring 条消息的环形缓冲区,当通过 /sub?last=N 连接时,先重新发送 seq 大于 N 的消息,再接着实时发送。如果 N+1 已经被挤出缓冲区,就先发送 {"reset": true, "seq": 最后的 seq}。然后在 /root/rt/wsops/client.py 中添加 consume(url, out_path, stop_after, base=0.1, cap=1.0) 协程。连接到 url,把收到的消息的 seq 每行一个地追加到 out_path 中,连接断开后,用 should_reconnect 和 backoff 等待一段时间,再把最后记录的 seq 作为 last 传入并重新连接。当 seq 到达 stop_after 时结束。评分器会在发布过程中两次断开连接,并查看文件中从 1 到 stop_after 是否不遗漏地各出现一次。
参考
- 工作文件夹是 /root/rt/wsops。请先用 mkdir -p /root/rt/wsops 创建。
- 想自己启动中心的话,请使用 cd /root/rt/wsops && /opt/rt-lab/bin/python -c "import asyncio, hub; asyncio.run(hub.run('127.0.0.1', 9002))" 这条命令。评分器会另行选择一个空闲端口来启动。
- 空闲超时和连接断开,通过 /opt/fixtures/rt/rtnet.py 的 IdleProxy 和 Relay.cut() 来制造。也可以读一读它是如何断开的。
- 有两个常见错误:在发布循环中对每个订阅者调用 await send,导致一个人让整体停住;以及在收到之后、记录 last 之前就断开,从而丢失消息。
- Python 必须用 /opt/rt-lab/bin/python 运行。本实验的库只装在那个虚拟环境里,如果直接用 python3 运行,就会出现 ModuleNotFoundError。像 alias rpy=/opt/rt-lab/bin/python 这样简写一下会比较方便。
- 实验 Pod 的对外连接被封锁。所有通信都发生在同一个 Pod 内的 127.0.0.1 上,不需要安装或下载。
- 评分器会以单独的进程加载你的代码,并实际建立连接。示例文件只是函数框架,原样保留是通不过的。前面步骤中已经完成的函数不要删除。
- 实验会话结束后,/root 中的文件不会保留。需要的代码请在结束之前另行保存。
给重连间隔掺入抖动
在 /root/rt/wsops/client.py 中创建 backoff(attempt, base=0.2, cap=5.0, rng=random)。返回第 attempt 次重试之前要等待的秒数。该值只通过一次 rng.uniform(0, min(cap, base * 2 ** attempt)) 来确定(full jitter)。rng 会传入 random 模块或 random.Random 对象。
如果没有抖动,一起断开的几万个客户端就会在完全相同的时刻重新涌来,把刚刚恢复的服务器再次压垮。如果没有上限(cap),第十次重试就会在几分钟之后。评分器会传入指定了种子的 random.Random 来计算相同的值。
根据关闭码决定是否重新连接
在 /root/rt/wsops/client.py 中添加 should_reconnect(code)。如果是 1001(离开)、1006(没有 close 帧就断开)、1011(服务器内部错误)、1012(服务重启)、1013(稍后重试),就返回 True,其他所有码都返回 False。
对于因 1008(违反策略)或 1002(协议错误)而断开的连接,重新连接也会因同样的原因再次断开。只有对方的情况有可能改变时,重试才有意义。1000 表示有人有意关闭了它。
让中心把一个人说的话传给所有人
在 /root/rt/wsops/hub.py 中创建 run(host, port, ring=1000, queue=100, ping_interval=1.0, ping_timeout=2.0) 协程。用 websockets.asyncio.server.serve 接收三个路径。对于通过 /pub 进来的每条文本消息,加上从 1 开始递增的 seq,做成 {"seq": n, "data": 消息} 形式的 JSON 字符串,发送给通过 /sub 连上的所有订阅者。/stats 发送一个 {"subscribers": 订阅者数, "dropped": 被断开的订阅者数, "seq": 最后的 seq} JSON 之后结束。为每个订阅者单独设置一个 asyncio.Queue,以及一个清空该队列并发送的循环。
如果发布一侧对每个订阅者依次调用 await send,那么一个人变慢时,排在他后面的所有人都要等待。每个订阅者一个队列,就是把这种等待按人分隔开的装置。发布只是放入队列,不做等待。
只断开一个慢订阅者
用 queue 限制订阅者队列的大小。如果队列已满、无法放入,就把该订阅者从列表中移除,dropped 加 1,然后用关闭码 1013 和原因 "slow consumer" 关闭。评分器会在接入一个不读取数据的订阅者和一个读取良好的订阅者的情况下,发布 8000 条 8000 字节的消息,查看读取良好的一方是否全部收到、中心的内存增长了多少,以及原来不读取的一方是否收到了 1013。
如果队列没有上限,一个慢订阅者的份额的消息就会在中心内存中无限堆积。丢弃消息和断开连接哪个更好,取决于数据的性质。对于顺序和遗漏都很重要的流,与其带着缺口继续发送,不如断开并让它重新连接,更为诚实。
不让安静的连接被负载均衡器切断
把 ping_interval 和 ping_timeout 按 run 的参数原样传给 serve。评分器会在中心前面放置一个会断开安静 2.5 秒的连接的中继器,用一个不会自行发送 ping 的订阅者等待 6 秒之后,发布一条消息。该消息必须到达订阅者。
websockets 默认的 ping 间隔是 20 秒,所以在空闲上限比它更短的设备后面,安静的连接会先被断开。常见的 AWS ALB 的默认空闲上限是 60 秒,而公司内部的代理往往更短。间隔必须比最短的空闲上限更短。
清除悄无声息消失的订阅者
修改订阅者循环,使它不只是等待队列,同时也等待连接关闭。用 asyncio.wait(..., return_when=FIRST_COMPLETED) 同时等待 ws.wait_closed() 和 q.get(),如果连接先关闭,就把它从订阅者列表中移除并结束循环。也要把 close_timeout=1.0 传给 serve。评分器会接入一个只做握手、不应答 ping 的订阅者,然后查看 6 秒内 /stats 的 subscribers 是否回到 0。
即使库因 ping 超时而关闭了连接,等待队列的协程在新消息到来之前也不会醒来。这期间,订阅者列表中仍留着已死的连接,使数字虚高,并且还在接收并堆积消息。而且,向不应答的对端发送 close 之后,等待应答的时间默认值是 10 秒。
从断开的位置接续
在中心中设置一个存放最近 ring 条消息的环形缓冲区,当通过 /sub?last=N 连接时,先重新发送 seq 大于 N 的消息,再接着实时发送。如果 N+1 已经被挤出缓冲区,就先发送 {"reset": true, "seq": 最后的 seq}。然后在 /root/rt/wsops/client.py 中添加 consume(url, out_path, stop_after, base=0.1, cap=1.0) 协程。连接到 url,把收到的消息的 seq 每行一个地追加到 out_path 中,连接断开后,用 should_reconnect 和 backoff 等待一段时间,再把最后记录的 seq 作为 last 传入并重新连接。当 seq 到达 stop_after 时结束。评分器会在发布过程中两次断开连接,并查看文件中从 1 到 stop_after 是否不遗漏地各出现一次。
必须以“处理并记录下来的”,而不是“收到的”为基准来重新连接。如果只是收到而在记录之前就断开,那条消息就永远丢失了。如果同一个 seq 来了两次,只记录一次。