The queue that never stops growing#
TL;DRthe 30-second version
- A queue absorbs short bursts, but it cannot raise the service rate. Under sustained overload the queue only grows, and the wait grows with it.
- The fix is a bounded queue: cap how many requests may wait, and reject new arrivals the instant it is full. That rejection is load shedding.
- Shedding keeps the accepted work fast. The excess is turned away in microseconds instead of sitting in a queue for seconds.
- Backpressure is the same idea aimed upstream: signal the sender to slow down instead of silently swallowing more than you can handle.
Picture a service with three workers. Each finishes one request per tick, so the service finishes three per tick. That is its service rate, and nothing you do to the queue changes it.
Requests wait in a queue until a worker is free. Now the load rises to five requests per tick and stays there. Five arrive, three leave. Every tick, two more requests pile up that never get worked off. After a hundred ticks the queue is two hundred deep. The workers are as busy as they can be, yet the backlog grows without limit.
The queue depth is not the part that hurts. The wait is. A request that lands when the queue is two hundred deep sits behind two hundred others, so it waits about sixty-six ticks before a worker even looks at it. By then the user has refreshed the page, and the retry lands at the back of an even longer queue. The service is now doing work whose answers nobody is waiting for. That is how an overloaded service topples.
Bound the queue, shed the rest#
The first move is to put a limit on the queue. Say the queue may hold at most six requests. When the seventh request arrives and the queue is full, the service does not buffer it. It rejects it immediately, before doing any real work, and moves on. Rejecting excess work at the door under overload is load shedding.
With the queue capped at six and three workers draining it, a request never sits behind more than six others. So it waits at most about two ticks, no matter how hard the overload pushes. A fast rejection is a far kinder answer than a response that arrives a minute too late. The caller learns immediately and can back off or try another replica.
Backpressure is the same principle pointed in the other direction. Instead of dropping what you cannot handle, you send a signal back up the chain: slow down. TCP does this at the network layer. The receiver advertises how much buffer space is left, and the sender is not allowed to send past it. Nothing has to be thrown away at all.
PredictUnder a steady overload of five arrivals per tick against three workers, you double the queue capacity from six to twelve. What happens to the worst wait an accepted request suffers?
Hint: The workers still finish three per tick no matter how long the line is.
It roughly doubles. A bigger queue does not add workers, so the drain rate is still three per tick. A request can now sit behind up to twelve others instead of six, so its wait grows from about two ticks to about four. The bigger buffer makes every accepted request slower. The service rate is the only real lever.
Two more tools decide which held work is worth doing. A deadline: if a request has waited longer than its caller will, drop it unserved and free the slot. And the serve order: under overload the oldest request is the one most likely to have already blown its deadline, so Facebook's Fail at Scale flips from first-in-first-out to last-in-first-out once the queue's wait time crosses a threshold. The newest requests are the ones that can still finish in time.
The numbers that decide it#
Little's Law says the number of items waiting in a system equals the arrival rate multiplied by the time each item spends there. Rearranged, the wait equals the number waiting divided by the rate they leave. Forty requests in the queue, three cleared per tick, is a wait of about thirteen ticks. Cap the number waiting and you have capped the wait. The same law sets the cap: the longest wait you will serve, times the service rate, is your queue size.
| Condition | What the queue does | What the wait does |
|---|---|---|
| arrivals ≤ service | stays shallow; absorbs short bursts and drains | flat and small |
| arrivals > service, no cap | grows by (arrivals − service) every tick, forever | grows without limit |
| arrivals > service, capped at C | fills to C and holds; the excess is shed | flat, at most about C ÷ service |
Choosing how to fail#
| Option | What happens | The cost |
|---|---|---|
| Shed | Reject the excess fast. Accepted work stays quick; rejected callers get an immediate, honest error. | Some requests are refused outright. |
| Buffer | Accept everything into a growing queue. Nothing is refused at the door. | Every request slows down, and past a point the whole service becomes unresponsive. Looks kindest, fails hardest. |
| Degrade | Serve a cheaper, worse answer to everyone: a cached subset, no personalized ranking. | Quality, in exchange for more throughput. |
| Push back | Signal the sender to slow down, so the excess is never generated. | Only works when you control the sender: a batch job, an internal client. Not open internet traffic. |
Real systems layer these. Google's SRE book describes the progression: push back on clients you control, degrade as you approach the ceiling, shed the excess, and serve outright errors only when even the degraded path is saturated. Netflix skips hand-tuning the cap. Its concurrency-limits library watches request latency, and when latency rises it lowers the limit automatically. The ideal limit is roughly the request rate times the minimum latency, which is Little's Law again.
If this comes up in an interview#
Isn't dropping requests just failing?
A fast rejection under overload is a better outcome than a slow success that arrives after the caller gave up, and far better than the whole service going down. Shedding fails a few requests quickly so the rest succeed quickly.
Backpressure or load shedding, which do I use?
Backpressure when you control the sender and can make it slow down: internal clients, batch jobs, a stream you consume. Load shedding when you cannot, like open internet traffic that keeps arriving no matter what you signal. Most systems use both, at different edges.
How big should the queue be?
Small. Size it from the longest wait you are willing to serve: that wait times the service rate, via Little's Law. A large queue does not help throughput. It only adds latency before you finally shed.
How is this different from rate limiting, a circuit breaker, or autoscaling?
Rate limiting is a per-caller budget decided ahead of time. A circuit breaker guards the calls you make to a failing dependency. Autoscaling raises the service rate over seconds or minutes. Backpressure is the reaction to the queue in front of you right now, with the capacity you have.
Common mistakes
- An unbounded queue anywhere. Every in-memory queue, thread pool, and channel needs a bound. One unbounded buffer in the chain is the most common cause of a service that runs out of memory under load instead of shedding cleanly.
- Shedding late. Rejecting a request after you have parsed, authenticated, and routed it burns the capacity you are protecting. Shed at the edge.
- Retries without backoff. A client that immediately retries a shed request re-adds the load you shed. Google's SRE book traces the cascade: a server tips into overload, client retries pile more load onto it, and the failure spreads to healthy machines. Pair shedding with exponential backoff and jitter.
- Serving stale work under overload. Finishing a request whose caller already timed out is wasted effort. Without deadlines, a backlog of stale requests keeps your workers busy producing answers nobody will read.
References
- Google SRE Book — Handling Overload — Graceful degradation, load shedding, and client-side throttling as the response to overload.
- Netflix — Performance Under Load (adaptive concurrency limits) — A TCP-congestion-control approach that infers overload from latency and sheds automatically.
- Facebook — Fail at Scale (controlled delay + adaptive LIFO) — Queue-sojourn-based shedding and serving newest-first under sustained overload.
- Little's Law — The relation between items in a system, arrival rate, and time spent — the math behind bounding the wait.