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

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

直播为什么会积压

在 TT Lab 中继续学习

一句话总结

队列会暂时保管待处理的工作,但不会让慢速的消费者变快。

为什么需要它

想象一下直播猫咪救援队的位置。九个订阅者会立即把新坐标画到屏幕上,但有一个订阅者的屏幕处理停住了。如果服务器依次等待所有订阅者完成,这一个停住的人会让其他九个人的地图也一起静止。反过来,如果先全部放进队列再去做下一件事,发布者是轻松了,但慢订阅者的队列会不断变大。这并没有消除等待,只是把等待转移到了内存里。

本课程面向了解 Python 函数、类、异常和 async/await 基础的学习者。如果你已学完前面的《TCP 包裹碎成几段才送到》课程,请回想一下:接收到字节与消息完成是两回事。这次要处理的问题是:即使消息已经完整,业务也可能积压。不需要攻击真实的外部服务或开放互联网,而是用同一个实验容器内的两个 TCP 订阅者来重现。

工作原理

每秒进入 100 个、处理 60 个的话,不丢弃的队列在处理速率不变期间,每秒大约增加 40 个。把队列扩大到 4,000 格,只是推迟了饱和的时间点,并不能解决持续的处理能力不足。不仅要看队列的长度,还要看等待最久的条目的年龄。100 个坐标是 1 秒的量还是 10 分钟的量,用户看到的画面含义是不同的。

asyncio.Queue(maxsize=4) 把队列中等待的条目限制为四个。worker 用 get 取走的那一个条目会从该长度中扣除,但在等待 ACK 期间仍然占用内存。每个连接有 4 个等待项和 1 个处理中的条目,此外还有序列化缓冲区和套接字缓冲区。因此不能把队列长度 4 称为整体内存 4 个条目的保证。还要限制条目的大小,才能算出以字节为单位的预算。

下表中,发布者是生产者,更新画面的 worker 是消费者。

调用 含义 常见误解
await queue.put(item) 等到出现空位再放入 所有订阅者都处理完了
queue.put_nowait(item) 立即放入,否则抛出 QueueFull 满了也会安静地等待
await queue.get() 取得一个等待条目的所有权 该条目的业务成功了
queue.task_done() 记录已取走工作的收尾 远程数据被永久保存了

事件循环在 await 让出控制权期间推进其他工作。没有 await 而无限循环的函数,会让其他连接也停住。async def 这个声明本身,并不会把函数变成并行的 CPU 工作。另外,asyncio.Queue 的 maxsize=0 不是零格,而是无限制。这就是第一步只接受正整数的原因。bool 在 Python 中是 int 的子类型,但作为这个 API 的容量会被拒绝。

如果取消在已满的队列中等待的 put,对应的发布也应该结束。如果捕获取消并像成功一样返回,调用方可能会误以为发送了其实没有放进去的事件。不要不经意地吞掉 CancelledError,要同时确认原有条目是否保持原样、被取消的条目日后是否不会出现。

在现场相遇的样子

实时字幕的生成器可能会暂时加快,小队列可以吸收瞬时的差距。但如果移动端屏幕一直很慢,就必须在某个时刻在等待、舍弃、断开连接中选择一个。订单处理不能丢弃旧订单,所以不使用与屏幕坐标相同的策略。在选择队列的实现之前,先问清数据的含义。

本课中的 ACK 搁置,是订阅者业务处理很慢的情形。它不是直接测量内核发送缓冲区饱和或互联网丢包的实验。必须分清测量的是哪一层,才不会误以为修改了队列长度就解决了网络故障。

下一项检查要做什么

请写出:往空的容量为 2 的队列里放入 A 和 B,worker 取走了 A。qsize 是 1,但因为还没有调用 task_done,未完成的工作有两个。放入 C 后 qsize 又变成 2,想放入 D 的生产者会等待。worker 取走 B 的那一刻出现空位,D 就可以放入。A 的业务完成与下一个条目可以入队的时刻,并不是同一件事。

在这里,如果 A 永远得不到 ACK,worker 就取不走 B。队列里留着 B 和 C,生产者停在 D。只看这个连接的队列,这是正常的背压。但如果这个生产者同时负责其他所有订阅者,就会蔓延成整个直播的停滞。即使是同一个工具,所有权范围和等待的位置不同,故障半径也不同。不要只看一个函数,请用箭头画出是谁在等待这个函数。

用测验来区分队列上限、处理中的条目和被取消的 put 的含义。也请追踪上面例子中 A 的 ACK 返回的情形,比较被阻塞的生产者以怎样的顺序重新推进。在最后一个模块的综合实验中,亲自实现 make_queue 和 offer_wait 之后,对同一个队列应用三种策略。