Backpressure and capacity
When work arrives faster than it can be done, something has to give: the queue grows, the producer slows, or work is shed. Choose which on purpose.
Performance & scale
Learn it
Every component has a maximum throughput. When arrivals exceed it, the excess accumulates somewhere: in a queue, in buffers, in open connections, in threads waiting on locks. It's invisible until latency explodes or memory runs out, and then the whole system fails, not just the excess.
Two pieces of arithmetic explain most capacity problems:
- Little's law: items in the system = arrival rate × time each spends there (L = λW).
- Utilization: as a server nears 100% busy, queueing delay grows roughly as 1 ÷ (1 − utilization). At 50% a request waits about one service time; at 90% about nine; at 99% about ninety.
Work it out
Requests arrive at 200 a second and each takes 50 ms. How many are in flight on average?When arrivals exceed capacity, there are only three responses:
- Buffer the excess: overload becomes latency. Fine for short bursts that drain in time; disastrous otherwise.
- Backpressure: bounded buffers that block or reject push the problem upstream, toward something that can decide (a client that retries later, a user who sees "busy"). TCP flow control is backpressure.
- Shed load: reject or drop work early so what you accept still finishes in time. See Load shedding.
Check
A service puts an unbounded in-memory queue in front of a slow dependency. Under sustained overload, what happens?
Quick reference
The same ideas, condensed for revision.
How it goes wrong
- Unbounded queue
- Latency grows without limit; by the time work is processed, nobody wants it.
- Slow consumer, fast producer
- Per-connection send buffers grow until the server runs out of memory.
- Retry amplification
- Rejected work is retried immediately, increasing the very load that caused rejection.
- Running hot
- Provisioning for 95% utilization means small bursts cause large latency spikes.
Instead, consider
- Add capacity (autoscaling)
- Load grows predictably or slowly enough for new capacity to arrive in time.
- Reduce work per request
- Caching, batching or coarser updates cut the service time itself.
In practice
- Bounded channels and queues
- Block or reject when full.
- HTTP 429/503 with Retry-After
- Explicit load shedding with guidance for clients.
- TCP and HTTP/2 flow control
- Receivers advertise how much they can accept.
- Per-connection buffer limits
- Disconnect or degrade clients that fall too far behind.
It assumes
- You can measure arrival rate, service time and queue age.
- Upstream callers can handle rejection or slow-down signals sensibly.
Explain it in your own words
Where you practise it
- A URL shortener like bit.ly
- A reliable video processing pipeline
Read the requirements like an engineer · 100x traffic: find the real bottleneck · Paid instructors want a fast lane
- Rate limiting a public API
What the numbers say · Redis goes down · Not all requests cost the same
- A home timeline at 300,000 reads a second
- A job queue that keeps working when workers fall behind
What the numbers say · Why the queue stopped draining · A buffer that can hold a bad day · Twenty million jobs waiting · Defend keeping Redis
- Product analytics over billions of events
- Notifications across email, push and in-app
What the numbers say · The email provider is failing · The announcement and the security alert
- A look-aside cache at Facebook's scale
The key everyone wants · A cache server dies · A cluster with an empty cache
- A real-time collaborative editor
- Sharding Postgres while it is running
Further reading
Engineers describing it in systems they run.
- Scaling Slack's Job Queue
Slack · Saroj Yadav and others · Post, Dec 2017
An outage, the design flaw behind it, and the redesign. The section on rolling the new path out without an outage is worth reading twice.
Related concepts
- Message queues
A durable buffer between producers and consumers that hands each message to one consumer at a time and redelivers it unless acknowledged.
- Rate limiting
Capping how fast a client may use a resource, to protect capacity, enforce fairness, and stay within the limits of the systems you depend on.
- Retries, backoff and jitter
Retrying transient failures with growing, randomized delays and a budget, so recovery does not become the next outage.
- Asynchronous processing
Accepting a request, recording the work durably, and doing it later in another process, so the work can outlive the request.
- Persistent connections
Long-lived connections such as WebSockets turn a stateless request tier into one that holds per-client state, with consequences for routing, deploys and failure detection.
- Partitioning
Splitting data or work by key so each part is handled independently: scaling out, and giving each key a single owner.
- Caching
Keeping a copy of data closer to where it is used, trading freshness and complexity for speed and reduced load on the source.
- Request coalescing
When many callers ask for the same thing at the same moment, do the expensive work once and give every caller the result. The fix for thundering herds and hot keys.
- Load shedding
Rejecting some work on purpose when a system is overloaded, so the work it does accept still finishes in time. Cheap rejections beat slow failures.