TT Lab
Get started
Learn Learning paths Courses

One Slow Subscriber Froze the Livestream

Why livestreams fall behind

Continue in TT Lab

In one line

A queue holds work for a while, but it does not make a slow consumer fast.

Why this was needed

Imagine you are streaming the location of a cat rescue team live. Nine subscribers draw the new coordinates on screen right away, but one has stopped processing its screen. If the server waits for every subscriber to finish in turn, the one stalled subscriber freezes the maps of the other nine as well. Conversely, if you put everything in a queue first and move on, the publisher is comfortable, but the slow subscriber's queue keeps growing. You did not remove the waiting; you moved it into memory.

This course is for learners who know Python functions, classes, exceptions, and the basics of async/await. If you finished the earlier TCP parcel course, recall that receiving bytes and completing a message are different things. This time we deal with the problem that work can back up even after a message is complete. You do not need to attack any real external service or open the internet; we reproduce it with two TCP subscribers inside the same lab container.

How it works

If 100 items arrive per second and 60 are processed, a queue that does not drop anything grows by about 40 items per second as long as the processing rate stays the same. Enlarging the queue to 4,000 slots only delays the moment it saturates and does not fix a sustained shortfall. You have to look not only at the length of the queue but also at the age of the item that has waited longest. Whether 100 coordinates are one second's worth or ten minutes' worth changes what the screen means to the user.

asyncio.Queue(maxsize=4) limits the items waiting in the queue to four. An item a worker takes with get leaves that length, but it is still in memory while it waits for the ACK. There are 4 waiting per connection and 1 in flight, and the serialization buffer and the socket buffer are separate on top of that. So you do not call a queue length of 4 a guarantee of 4 items of total memory. You also have to limit the size of each item before you can set a byte-level budget.

In the table below, the publisher is the producer and the worker that updates the screen is the consumer.

Call Meaning Common misunderstanding
await queue.put(item) Waits until a slot opens, then inserts Every subscriber has finished processing
queue.put_nowait(item) Inserts right away or raises QueueFull It quietly waits even when full
await queue.get() Takes ownership of one waiting item The work for that item succeeded
queue.task_done() Records cleanup of a taken task The remote data was stored permanently

The event loop runs other tasks while you yield control with await. A function that loops forever without await also stalls the other connections. The async def declaration itself does not turn a function into parallel CPU work. Also, maxsize=0 in asyncio.Queue means unlimited, not zero slots. That is why the first step accepts only positive integers. In Python, bool is a subtype of int, but this API rejects it as a capacity.

If you cancel a put that is waiting on a full queue, that publish must end as well. If you catch the cancellation and return as if it succeeded, the caller may believe it sent an event that never went in. Do not casually swallow CancelledError, and check together that the existing items remain as they were and that the canceled item does not show up later.

What it looks like in the field

Live captions can have a generator that briefly speeds up, so a small queue absorbs the momentary difference. But if the mobile screen stays slow, at some point you have to choose one of waiting, skipping, or closing the connection. Order processing cannot drop old orders, so it does not use the same policy as screen coordinates. Before you choose a queue implementation, the first thing to ask is what the data means.

The ACK hold in this lesson is a situation where the subscriber's work processing is slow. It is not an experiment that directly measured kernel send buffer saturation or internet packet loss. You have to be clear about which layer you measured so that you do not mistake changing the queue length for having solved a network failure too.

What you will do in the next check

Write down what happens when you put A and B into an empty queue of capacity 2 and the worker takes A. qsize is 1, but since task_done has not been called yet, there are two unfinished tasks. If you put in C, qsize is 2 again, and a producer trying to put in D waits. The moment the worker takes B, a slot opens and D can go in. The completion of A's work and the moment the next item can be inserted into the queue are not the same event.

Here, if A is never ACKed, the worker cannot take B. B and C stay in the queue and the producer stops at D. Looking only at this connection's queue, that is normal backpressure. But if that producer is responsible for every other subscriber, it spreads into a halt of the whole livestream. Even the same tool has a different blast radius depending on its ownership scope and where it waits. Do not look at a single function alone; draw with arrows who is waiting on that function.

You distinguish the queue limit, in-flight items, and the meaning of a canceled put in the quiz. Also trace the case where A's ACK comes back in the example above and compare in what order the blocked producer makes progress again. In the integrated lab in the last module, you implement make_queue and offer_wait yourself and then apply three policies to the same queue.