明确确认、期限与清理的责任
一句话总结
发送缓冲区已空、确认业务已完成、连接已回收,这是三种不同的证据。
为什么需要它
服务器完成 write 之后记录了成功日志,用户的屏幕却没有变化。因为字节被交给操作系统,与对方程序完成业务,是不同的。反过来,对方可能已经处理了,但在 ACK 到来之前连接断了。实时程序不能隐藏这种不确定性,要用消息契约、期限和恢复位置来表达。
这次的契约非常小。服务器以 ASCII 发送 EVENT 7 加换行,客户端假定已完成该业务,然后返回 ACK 7 加换行。序号是从 0 到 2147483647 的整数。它不是经过认证的业务协议,也不是持久化的消息代理,而是用来学习确认和生命周期的教学用协议。
工作原理
StreamWriter.write 把数据放进写缓冲区。await drain 会等到传输流控允许为止。drain 的完成并不是对方应用的完成,所以接着用 reader.readline 等待同一序号的 ACK。EOF、其他序号、错误格式都按 ConnectionError 处理。发送函数不会随意关闭被多个事件使用的持久连接。
时间预算对每个事件分配一次。总共 1 秒,如果 drain 用了 0.8 秒,ACK 就只剩下约 0.2 秒。如果给每个 await 都重新分配 1 秒,最长等待时间就会随任务数增加。开始时用 monotonic 时钟确定截止时刻,并在每次等待之前计算剩余时间。最后一个结果回来之后,也要确认没有超过期限。不要用墙上时钟的校准来判定经过时间。
队列 worker 用 get 接管一个条目之后执行 send_one。无论成功还是异常,接管条目的收尾都要在 finally 中用 task_done 完成。如果在 ACK 之前调用 task_done,发布一侧的 join 就会比实际确认更早解除。反过来,如果异常时漏掉 task_done,工作已经死了,join 还在一直等待。在这里再次应用:收尾与业务成功是不同的记录。
如果 worker 函数接管了连接的所有权,那么无论在空队列上被取消,还是在等待 ACK 时失败,都必须调用 writer.close。也要给 wait_closed 设置有限的期限,如果终止没有完成,就用 transport.abort 强制回收。队列中仍然剩下的条目的丢弃和重试,是订阅者管理器的责任。不要把单个 worker 清理自己处理中的条目,与清空整个队列混为一谈。
重新连接需要最后一个确定处理过的序号。Cursor(last=7) 收到 8 是连续推进,再次收到 7 则是重复。如果 10 先到,8 和 9 就是空缺,所以抛出 Gap,并且不移动 last。要先通过日志重放或可信的快照恢复缺失,再继续推进。如果直接采用较大的序号,缺失就会被永久隐藏。
这里构造的 Cursor,是在内存中判断连续性的小部件。accept 为 True 时会立即推进 last,所以不能在调用可能失败的真实支付函数之前直接使用。本课程没有以原子方式保存业务应用和检查点的功能。进程重启时,内存也会消失。通过了重复判断的测试,并不保证恰好一次处理。
在现场相遇的样子
在生产中,订阅者连接断开、发布者取消、进程终止会在不同的时刻叠加。如果只有正常 ACK 的测试成功,就会漏掉这些终止路径中的泄漏。要确认即使反复断开连接,也不会遗留套接字和任务。先关闭监听器阻止新连接,再清理现有任务和连接,最后等待服务器终止,这个顺序也很重要。Python 3.12 的 Server.wait_closed 会等待活动连接终止。
最后的真实 TCP 测试观察两个订阅者的部分故障和回收。前面的步骤用假时钟确定性地检查整体期限,ACK 永远不来的情形也用真实的异步等待来检查。两种方式不是替代关系。假时钟确认边界计算,真实连接确认接线和生命周期。
下一项实验要做什么
在期限测试中,不必真的多次等待 1 秒,可以拨动注入的时钟来检查边界。但如果假的 reader 立即返回,即使没有中断无限等待的代码,也可能只有那个例子通过。所以还要另外使用永远不返回的 reader。如果因为失败测试耗时长就把它去掉,那个缺陷在生产中就会变成无限等待。这就是为什么要同时保留计算边界和真实取消。
在第 8 步中完成队列策略、游标、ACK 整体期限、worker 和 broadcast。所有代码放在 /root/realtime/delivery.py 中,服务器由检查器通过临时端口准备。在查看正确答案之前,请先读一读失败消息要求的是哪一层的证据。