An ingestion service accepts events over HTTP and writes them to a database. It has a rate limit of 5,000 requests per second, chosen a year ago by measuring what the service could handle and adding some margin.
The database gets slower one afternoon because a large migration is running. Write latency goes from 4 milliseconds to 90.
The rate limiter continues to admit 5,000 requests per second, because that is the number it was given and nothing about the database is visible to it. The internal queue between the HTTP layer and the database writer grows. Memory climbs. Latency goes from 20 milliseconds to 6 seconds, then to 40. The pod is eventually killed for memory and restarts, losing everything queued.
The rate limiter did its job perfectly. Its job was never this.
Two mechanisms, two different failures
Rate limiting answers: is this caller allowed to send this much? It is a policy decision made in advance, and the number comes from a contract, a business tier, or a measurement taken at some point in the past.
Backpressure answers: can this system currently absorb more? That is a runtime question with a runtime answer, and the answer changes when a dependency slows down, a node is lost, or a noisy neighbour takes CPU.
flowchart LR
subgraph rl["Rate limiting"]
P1["Producer"] -->|"capped at a fixed number"| S1["Service"]
S1 -.->|"no signal back"| P1
end
subgraph bp["Backpressure"]
P2["Producer"] -->|"send"| S2["Service"]
S2 -->|"slow down, I am full"| P2
end
style S1 stroke:#ef4444,stroke-width:2px,color:#fff
style S2 stroke:#4ade80,stroke-width:2px,color:#fff
The direction of the arrow is the whole distinction. Rate limiting is a constraint applied at the entrance. Backpressure is information flowing backwards from wherever the actual constraint currently is.
Both are worth having. A rate limit stops one tenant consuming everything and gives you a defensible answer about fairness. It cannot protect you from your own capacity dropping, because it has no idea that happened.
The unbounded queue is what makes it fatal
Between any producer and consumer there is a buffer, and whether it has a limit determines whether overload is survivable.
// Unbounded. Accepts work forever. The failure is memory, later, with no warning.
ExecutorService bad = new ThreadPoolExecutor(
10, 10, 0L, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<>()); // no capacity argument
// Bounded, with an explicit decision about what happens when it is full.
ExecutorService good = new ThreadPoolExecutor(
10, 10, 0L, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<>(1_000),
new ThreadPoolExecutor.AbortPolicy()); // reject, loudly and immediately
Executors.newFixedThreadPool uses an unbounded queue. So does newSingleThreadExecutor. Both are the first thing most people reach for, and both will accept work until the process dies.
An unbounded queue does not remove the overload. It converts it from an immediate rejection, which the caller can handle, into unbounded latency followed by a memory failure, which nobody can handle. It also destroys the evidence, because the request rate looks fine and the latency graph is the only thing that moves.
Once the queue is bounded, the interesting decision is the rejection policy.
AbortPolicy throws and the caller learns immediately, which is usually what you want at a service boundary.
CallerRunsPolicy makes the submitting thread execute the task itself. That is backpressure in its purest form: the producer physically cannot submit more while it is busy doing the work, so the rate self-regulates with no coordination. It is elegant and it has a sharp edge, because if the caller is an HTTP thread you have just made a request handler run a background task.
DiscardOldestPolicy is right for telemetry and metrics where recent data matters more than complete data, and wrong for anything a user is waiting on.
The mistake is not picking one. It is having no bound, which means the policy never runs.
Propagating the signal across a boundary
Within a process, a bounded blocking queue propagates naturally: a producer that cannot enqueue blocks, and its caller experiences that as slowness, which travels up the chain.
Across a network there is no shared blocking, so the signal has to be explicit.
HTTP 429 Too Many Requests -> you exceeded your allowance (rate limit)
HTTP 503 Service Unavailable -> I cannot cope right now (backpressure)
Retry-After: 2 -> and here is when to come back
Those two status codes mean different things and are frequently used interchangeably, which matters because a well behaved client should respond to them differently. A 429 means slow down to your agreed rate. A 503 means the far side is in trouble and the correct response is to back off substantially, with jitter, and possibly to open a circuit.
In reactive streams the mechanism is inverted and explicit: the consumer requests a number of items, and the producer may not send more than were requested. Demand flows backwards by construction rather than being inferred from timing. gRPC streaming and HTTP/2 do the same at the transport layer with flow control windows, which is one of the quieter benefits of HTTP/2 multiplexing over a single connection.
For queue based systems, consumer lag is the backpressure signal, and it is worth treating as one rather than as an alert threshold. Lag growing steadily means the consumer cannot keep up, and the useful response is upstream: slow the producer, shed lower priority work, or scale the consumer.
Load shedding is backpressure with a decision attached
At some point the honest answer is to refuse work, and refusing well is a design problem.
Shedding at the entrance is much better than shedding in the middle, because a request rejected before any work is done costs almost nothing, while one rejected after three downstream calls has already consumed capacity that produced nothing.
Priority matters more than most systems implement. Rejecting a health check to serve a checkout is obviously correct, and it requires the system to know which is which. Without a notion of priority, shedding is random, and random shedding under load means every user gets a degraded experience rather than most users getting a working one.
The queue timeout is the piece that is almost always missing. An item that has been queued for 30 seconds is usually worthless, because the client gave up at 5. Processing it consumes capacity to produce a result nobody will read, which is why a deadline that travels with the request and is checked before execution recovers real capacity during overload.
public void submit(Task task) {
task.setDeadline(Instant.now().plusSeconds(5));
if (!queue.offer(task)) {
throw new OverloadedException(); // fail now, cheaply
}
}
public void execute(Task task) {
if (Instant.now().isAfter(task.getDeadline())) {
metrics.increment("task.expired");
return; // nobody is waiting for this
}
doWork(task);
}
What I would look for
Find every unbounded queue. Executors.newFixedThreadPool, LinkedBlockingQueue with no capacity, any channel or buffer created without a size. Each one is a place where overload becomes invisible until it becomes fatal.
Export queue depth as a metric everywhere there is a queue. It is the earliest signal of a capacity problem and it usually moves well before latency does.
Distinguish 429 from 503 in both directions, and make sure your clients treat them differently, because a client that retries a 503 the same way it retries a 429 is amplifying an outage rather than backing off from one.
The framing I keep returning to is that a rate limit encodes what you believed a year ago, and backpressure encodes what is true right now. Systems that only have the first one are protected against the failure they anticipated and open to the one they did not.
// SPONSORSHIP
If this research saved you time or improved your architecture, consider sponsoring my work on GitHub. All sponsorships go directly toward infrastructure and further technical research.
[ Become a Sponsor ]