One Slow Subscriber Froze the Livestream
What should we drop, and whom should we wait for?
In one line
The latest location and a payment history are not the same kind of event. The slow-subscriber policy changes depending on what you can afford to lose.
Why this was needed
The rescue team has already moved on to the next alley, yet the map replays its position from 5 minutes ago one step at a time. The server reports that it has not lost a single item, but it is a useless livestream to the user. Conversely, if you delete old entries from a payment history and show only the latest, you can drop the movement of money itself. Eliminating loss and serving the purpose of the service are not always the same choice.
This module compares three policies: waiting until space opens up, discarding the state that has waited longest, and separating a slow subscriber from the publish targets. All of them have pros and cons. Do not think that setting one queue size also decides the policy automatically. What you do after QueueFull determines the product's behavior.
How it works
The wait policy can be implemented with await queue.put. It does not quietly discard data, but if the publisher goes through the subscriber list and waits for each put, the whole pass stops at a single slow subscriber. Giving each subscriber its own queue does not isolate them automatically. You also have to look at in what order and in what way the producer waits on those queues.
The latest-state policy takes out the one item that has waited longest when the queue is full and puts in the new item. With capacity 2 and 10, 11 waiting, if 12 arrives, 11 and 12 remain. If you discard 12, which just arrived, then even though the name is latest, it is in practice a policy that preserves an old screen. An item that a worker has already taken and is processing cannot be removed from the queue, so this policy does not rewind work in progress.
A discarded item was also taken from the queue with get, so you match it with one task_done. Otherwise queue.join can wait forever. Even then, join finishing does not mean everything succeeded. Some were discarded on purpose and some may have failed to process. The counts of success, discard, and failure must be recorded separately as business metrics. In the lesson, we return the discarded item so you can compare the policies by eye.
The disconnect policy does not quietly modify the full queue but raises a SlowConsumer exception. The publish function removes that subscriber from the dictionary and asks the owner to terminate it. It keeps delivering to the remaining subscribers. If you return at the first failure, healthy subscribers later in the list miss the event. Iterating over a dictionary while it is being modified can cause an iteration error, so you use a snapshot of the name-queue pairs.
| Data | Policy to consider | What else you need |
|---|---|---|
| Current position, progress | Skip old waiting states | The latest snapshot, a marker for what was missed |
| Chat and notification history | Disconnect the slow connection, then reconnect | A retained log, the last confirmed position |
| Payment and inventory changes | Durable storage and retries | Idempotent handling, atomic state changes |
Disconnecting frees server resources but does not solve the delivery problem. If the subscriber received the event but was disconnected before sending the ACK, the server does not know whether it was processed. If you retry, duplicates are possible, and if you do not retry, omissions are possible. Disconnecting is a measure that reduces the scope of a failure and is a design separate from a delivery guarantee.
What it looks like in the field
A dashboard may not need to draw every intermediate value of a CPU utilization that changes dozens of times per second. But if you compress an audit log the same way, you lose the incident path. You have to state the policy per channel and observe together the number of times slow connections were disconnected and the time subscribers take to resynchronize. Even if only the disconnect count goes up, the latency metric of the healthy subscribers can look good.
The last step of this lab verifies the disconnect policy on real TCP. fast ACKs every event and slow holds back the ACK of the first event. When you publish from 0 to 8, you check that fast receives all of them and that only slow is disconnected, once. You do not lean on a fixed sleep; you check the execution order using ACKs and event barriers. This is not a throughput benchmark but a functional check of a partial failure.
What you will do in the next check
Write down the difference between the policies with the same input. slow is processing 0, 1 and 2 are waiting in a capacity 2 queue, and 3 arrives. The wait policy stops the publisher. The latest-state policy discards 1 and keeps 2 and 3. The disconnect policy rejects this insertion and notifies the connection owner. You must also decide who cleans up 1 and 2 in the disconnected queue. In no case is the success of 0, which is already in flight, settled by itself.
If you compress this difference into a single success or failure log line, you miss the cause. Queue saturation, a work ACK timeout, a remote EOF, and an administrator's forced disconnect are different events. If you distinguish metric names and failure messages, it is easier to judge whether you need to increase reconnects, improve the consumption speed, or change the data policy. Simply catching every error and continuing to send can produce a result where the livestream looks alive while only the data disappears.
In the quiz you check who each policy makes wait and what it loses. In the lab that follows, you implement offer_latest, offer_disconnect, and broadcast and explain why the same input is handled differently.