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

一个慢订阅者让直播停了下来

隔离慢订阅者,恢复直播传递

在 TT Lab 中继续学习

目标

实现队列策略、ACK 期限和连接生命周期,在真实 TCP 上隔离慢订阅者。

为什么重要

一个订阅者的延迟,可能会阻塞整个发布,或者让内存持续被占用。区分数据的可丢失范围与业务确认,并编写即使在取消时也会回收资源的程序。需要掌握 Python 函数、异常和 async/await 的基础。只使用标准库,不需要安装任何东西,也不需要互联网。

步骤

  1. 给无限堆积的队列设上限——在 delivery.py 中实现 make_queue(capacity)。返回一个空的 asyncio.Queue,其 maxsize 必须与输入相同。只允许正的 int,bool、浮点数、字符串、0 和负数都是 ValueError。
  2. 等待,但不留下被取消的发布——添加 async offer_wait(queue, item)。等到出现空位再放入,并返回 None。等待中被取消时,以 CancelledError 传递,并保留原有条目及其顺序。不要在日后插入被取消的条目。
  3. 丢弃旧坐标,保留当前位置——添加同步函数 offer_latest(queue, item)。有空位时插入后返回 None,队列已满时移除等待最久的一个条目,放入新条目,然后返回被移除的条目。对被丢弃的条目也要对应调用一次 task_done。
  4. 用策略异常通知已停滞的订阅者——实现 Exception 子类 SlowConsumer 和同步的 offer_disconnect(queue, item)。有空位时插入后返回 None。队列已满时不改变原有队列,而是抛出 SlowConsumer。关闭连接不是这个函数的职责。
  5. 跳过重复,不隐藏缺失——添加 Exception 子类 Gap 和 Cursor(last=-1)。last 是每个实例的公开属性,是范围为 -1–2147483647 的 int。accept(seq) 只接受 0–2147483647 的 int。seq<=last 时返回 False 并保持状态,seq==last+1 时更新 last 后返回 True,更大的序号则抛出 Gap 并保持状态。bool 等错误的类型和范围是 ValueError。它是内存游标,不实现持久化的业务处理。
  6. 对发送和业务确认只给一次时间——创建 async send_one(reader, writer, seq, timeout, clock=None)。seq 是与第 5 步相同的有效序号,timeout 是排除 bool 的正的有限 int/float,不合法则为 ValueError。默认时钟是 time.monotonic,注入的时钟不带参数并返回秒数。用 writer.write 一次性写入 ASCII 的 EVENT 序号加换行,等待 drain,然后用 reader.readline 收到序号完全相同的 ACK 加换行,才返回 None。EOF 和错误的 ACK 是 ConnectionError,整体期限超过是 TimeoutError。drain 和 ACK 共用同一个截止时刻,完成后也要检查期限。这个函数不关闭 writer,而是传递错误和取消。
  7. 即使队列为空或失败,也要回收连接——添加 async serve_queue(reader, writer, queue, timeout)。传入正的有限 timeout。循环地用 send_one 处理 get 到的 seq,只有在该调用结束、失败或被取消之后,才对接管的条目恰好调用一次 task_done。包括空队列等待在内,在所有结束路径上都要 writer.close 之后在 timeout 内等待 wait_closed。终止等待超时或出现 ConnectionError 时,用 writer.transport.abort 回收,并且不要隐藏原来的发送错误和取消。仍留在队列中的条目的清理由调用方负责。
  8. 分离慢的那一个,继续直播——完成同步的 broadcast(queues, item, disconnect)。queues 是名称到队列的字典,对遍历快照中的每个队列应用 offer_disconnect。只把抛出 SlowConsumer 的对象从字典中移除,调用一次 disconnect(name) 之后,继续向其他订阅者传递。其他异常要传递出去,正常返回 None。回调不抛出异常,负责取消对应的 worker、回收连接和清理剩余队列。在真实的 TCP 检查中,slow 搁置第一个 ACK,fast 每次都确认。事件 0–8 必须由 fast 全部收到,并且只有 slow 被分离。

参考

所有函数都放在 /root/realtime/delivery.py 一个文件中。请用 mkdir -p /root/realtime 创建工作目录。示例中尚未实现的函数保留为框架,但不要覆盖已完成的函数。评分也会检查之前步骤的契约。一个步骤的执行限制为 5 秒,这并不是对学习者编写时间的限制。TCP 服务器、临时端口和订阅者由检查器准备并终止。本实验不保证内核缓冲区饱和、互联网性能、持久传递和真实的重新连接恢复。实验会话结束后文件不会保留,请在结束前另行保存需要的代码。

给无限堆积的队列设上限

在 delivery.py 中实现 make_queue(capacity)。返回一个空的 asyncio.Queue,其 maxsize 必须与输入相同。只允许正的 int,bool、浮点数、字符串、0 和负数都是 ValueError。

在 asyncio.Queue 中 0 表示无限制。请把类型检查和范围检查分开,并使用标准库。

等待,但不留下被取消的发布

添加 async offer_wait(queue, item)。等到出现空位再放入,并返回 None。等待中被取消时,以 CancelledError 传递,并保留原有条目及其顺序。不要在日后插入被取消的条目。

put_nowait 和 await put 在饱和时的行为不同。不要把取消变成正常的成功。

丢弃旧坐标,保留当前位置

添加同步函数 offer_latest(queue, item)。有空位时插入后返回 None,队列已满时移除等待最久的一个条目,放入新条目,然后返回被移除的条目。对被丢弃的条目也要对应调用一次 task_done。

容量为 2 的队列中有 10、11 时,来了 12 就剩下 11、12。处理中的条目无法从这个队列中移除。

用策略异常通知已停滞的订阅者

实现 Exception 子类 SlowConsumer 和同步的 offer_disconnect(queue, item)。有空位时插入后返回 None。队列已满时不改变原有队列,而是抛出 SlowConsumer。关闭连接不是这个函数的职责。

函数负责判断饱和策略,连接所有者负责回收。悄悄丢弃或等待属于另外的策略。

跳过重复,不隐藏缺失

添加 Exception 子类 Gap 和 Cursor(last=-1)。last 是每个实例的公开属性,是范围为 -1–2147483647 的 int。accept(seq) 只接受 0–2147483647 的 int。seq<=last 时返回 False 并保持状态,seq==last+1 时更新 last 后返回 True,更大的序号则抛出 Gap 并保持状态。bool 等错误的类型和范围是 ValueError。它是内存游标,不实现持久化的业务处理。

重复检查与缺口检查的顺序很重要。7 之后来了 9,就没有处理过 8 的依据。

对发送和业务确认只给一次时间

创建 async send_one(reader, writer, seq, timeout, clock=None)。seq 是与第 5 步相同的有效序号,timeout 是排除 bool 的正的有限 int/float,不合法则为 ValueError。默认时钟是 time.monotonic,注入的时钟不带参数并返回秒数。用 writer.write 一次性写入 ASCII 的 EVENT 序号加换行,等待 drain,然后用 reader.readline 收到序号完全相同的 ACK 加换行,才返回 None。EOF 和错误的 ACK 是 ConnectionError,整体期限超过是 TimeoutError。drain 和 ACK 共用同一个截止时刻,完成后也要检查期限。这个函数不关闭 writer,而是传递错误和取消。

每次都要把剩余时间交给 asyncio.wait_for。测试会注入时钟,检查 drain 0.8 秒加上 ACK 0.3 秒是否超过整体的 1 秒。

即使队列为空或失败,也要回收连接

添加 async serve_queue(reader, writer, queue, timeout)。传入正的有限 timeout。循环地用 send_one 处理 get 到的 seq,只有在该调用结束、失败或被取消之后,才对接管的条目恰好调用一次 task_done。包括空队列等待在内,在所有结束路径上都要 writer.close 之后在 timeout 内等待 wait_closed。终止等待超时或出现 ConnectionError 时,用 writer.transport.abort 回收,并且不要隐藏原来的发送错误和取消。仍留在队列中的条目的清理由调用方负责。

条目所有权的 try/finally 与连接所有权的 try/finally,范围是不同的。如果在 ACK 之前调用 task_done,join 就会先被解除。

分离慢的那一个,继续直播

完成同步的 broadcast(queues, item, disconnect)。queues 是名称到队列的字典,对遍历快照中的每个队列应用 offer_disconnect。只把抛出 SlowConsumer 的对象从字典中移除,调用一次 disconnect(name) 之后,继续向其他订阅者传递。其他异常要传递出去,正常返回 None。回调不抛出异常,负责取消对应的 worker、回收连接和清理剩余队列。在真实的 TCP 检查中,slow 搁置第一个 ACK,fast 每次都确认。事件 0–8 必须由 fast 全部收到,并且只有 slow 被分离。

第一次分离之后不要 return。检查器会启动真实的服务器和两个客户端,所以不需要自己常驻运行服务器。