Short answers to what readers ask most about this topic.
01What is backpressure in simple terms?
Backpressure is a slow consumer pushing back on a fast producer so the producer slows down, pauses or drops work. It keeps the buffer between them from growing without limit. The signal can be a full bounded queue, a false return from write(), a TCP window, or an HTTP 429 or 503.
02How do you stop a fast producer from crashing a slow consumer?
Bound every queue and decide what happens when it is full: block the producer, reject the new work or drop the least valuable item. Let the consumer set the pace with pull-based consumption or a drain event. At the public edge, shed load with 429 or 503 and a Retry-After header.
03What does write() returning false mean in Node.js?
It means the stream wants you to wait for the drain event before writing more data. The Node.js documentation says not to write more chunks after false until drain is emitted. Ignoring it makes the internal buffer grow with every extra write, as the demo in this post shows.
04Should an overloaded API return 429 or 503?
Return 429 Too Many Requests when one client has exceeded its own limit, and 503 Service Unavailable when the whole service is overloaded. Both can carry a Retry-After header with a number of seconds or an HTTP date. Clients should add random jitter so rejected callers do not all retry at the same moment.
05Is an unbounded queue ever acceptable?
Only when you know the arrival rate cannot stay above the service rate for long, for example a small, finite batch you generate yourself. For anything driven by outside traffic, an unbounded queue just delays the failure and adds waiting time. Bound it and choose a full-queue policy instead.
Backpressure in Distributed Systems: Stop a Fast Producer
Backpressure is how a slow consumer tells a fast producer to slow down. Learn bounded queues, load shedding, Node.js drain and queue-depth autoscaling.
Backpressure is a signal from a slow consumer that tells a fast producer to slow down, so queues stay bounded instead of growing until memory runs out. Implement it with bounded queues, pull-based consumption, write() returning false then waiting for drain in Node.js, TCP windows, or load shedding with 429 or 503 and Retry-After.
A queue looks like free insurance: when the consumer is slow, the producer just keeps adding and the queue absorbs the difference. That only holds while the consumer is, on average, at least as fast as the producer. The moment it is not, the queue stops being a buffer and becomes a slow memory leak with a latency attached.
I run Docker, Redis and Postgres on a single VPS, so scaling out of a growing queue is not an option I can assume. This post answers one question: how do you stop a fast producer from crashing a slow consumer? It works through the arithmetic with Little's law, then bounded queues, load shedding, pull-based workers, Node.js streams and TCP, and ends with a queue-depth autoscaling rule. The Node.js demo output is real, from a run I made for this article.
What is backpressure in a distributed system?
Backpressure is any mechanism by which a consumer that cannot keep up pushes resistance back toward the producer, so the producer slows, pauses or drops work instead of piling it up. The word comes from fluid dynamics, where a restriction downstream raises pressure upstream. In software the restriction is a slow database, a rate-limited API or a busy worker pool.
The key idea is that the signal travels against the data flow. Without it, the producer only learns about trouble from a timeout or a crash, long after the damage is done. The Reactive Streams initiative states the goal plainly: governing the exchange of stream data across an asynchronous boundary so the receiver is not forced to buffer unlimited data.
What happens when a producer is faster than its consumer?
The queue grows without limit, memory climbs and waiting time grows with it. Little's law relates the average number of items in a stable system to the arrival rate and the time each item spends there. It needs a steady state, so for the overloaded case the arithmetic is simply arrival rate minus service rate. Both cases are worked below with assumed figures.
Stable system (Little's law, L = lambda x W)
lambda = 100 requests/s arriving, W = 0.05 s inside the service
L = 100 x 0.05 = 5 requests in the system at any moment
Same traffic, downstream now slow (W = 2 s)
L = 100 x 2 = 200 requests in flight
at an assumed 256 KiB of buffered state each: 200 x 256 KiB = 50 MiB
Unstable system (arrivals beat the service rate, so there is no steady state)
lambda = 200 messages/s in, mu = 150 messages/s out
queue grows by 200 - 150 = 50 messages/s
after 1 hour: 50 x 3600 = 180,000 messages queued
at an assumed 2 KiB each: 180,000 x 2048 = 368,640,000 bytes = 351.6 MiB
the next message waits 180,000 / 150 = 1,200 s = 20 minutes
Notice two separate costs. The memory cost is bounded by the machine, so it ends in an out-of-memory kill. The latency cost is bounded by nothing: a message that waits 20 minutes in the queue is often already useless to the user who sent it, so the consumer spends its capacity on work nobody is waiting for.
An unbounded queue does not remove the overload, it hides it. The failure moves from a fast, visible error at the edge to a slow crash somewhere you are not looking, usually after the queue has already ruined latency for everyone.
Is a bounded queue the fix?
A bounded queue is the minimum viable form of backpressure, because it forces a decision at the moment the queue is full instead of never. Bounding is not the fix by itself. It converts an unbounded memory problem into a question: what should the producer experience when the queue is full?
There are three honest answers, and each fits a different workload.
Block the producer: the push call waits for room. Good when the producer is another internal worker that can safely pause, as in a stream pipeline.
Reject the new work: the push call fails fast and the caller decides. Good at an edge that faces users or other services, and the basis of load shedding.
Drop the oldest or least valuable item: keep the freshest data. Good for metrics, positions and anything where a stale value is worthless, and wrong for payments or orders.
How does load shedding work with 429, 503 and Retry-After?
Load shedding refuses work at the edge once a concurrency limit is reached, so the work that is accepted still finishes quickly. HTTP gives you the vocabulary. RFC 9110 defines 503 Service Unavailable as the server being unable to handle the request due to temporary overload, and allows a Retry-After header, as a number of seconds or an HTTP date, to suggest how long the client should wait.
import { createServer } from "node:http";
const MAX_IN_FLIGHT = 50; // sized from Little's law: arrival rate x acceptable latency
let inFlight = 0;
createServer(async (req, res) => {
if (inFlight >= MAX_IN_FLIGHT) {
// Refuse in microseconds instead of queueing the work and timing out later.
res.writeHead(503, {
"Retry-After": "5", // delay-seconds form, RFC 9110 section 10.2.3
"Content-Type": "application/json",
});
res.end(JSON.stringify({ error: "overloaded", retryAfterSeconds: 5 }));
return;
}
inFlight++;
try {
await handle(req, res); // your real handler
} finally {
inFlight--; // always release, or the limit leaks shut
}
}).listen(3000);
// Per-client limit instead? Same shape, but answer 429 Too Many Requests.
Use 429 Too Many Requests when one client is over its own limit, and 503 when the whole service is saturated. RFC 6585 defines 429 as the user having sent too many requests in a given time and allows Retry-After there too. RFC 9110 also notes that a server does not have to use 503 when overloaded; some simply refuse the connection.
Retry-After only works if clients obey it. Make clients wait the stated seconds plus random jitter, otherwise every rejected caller returns at the same instant and recreates the spike you just shed.
How does Node.js handle backpressure in streams?
The Node.js stream documentation says writable.write() returns false if the stream wishes for the calling code to wait for the drain event before continuing to write more data, and that once it returns false you should not write more chunks until drain is emitted. The demo below ignores that rule in one loop and obeys it in the other. The consumer is deliberately slow, and the highWaterMark is set to 16 objects.
import { Writable } from "node:stream";
import { once } from "node:events";
const ITEMS = 100_000;
const PAYLOAD_BYTES = 1024;
const HIGH_WATER_MARK = 16; // objectMode: counts objects, not bytes
// A deliberately slow consumer: each item needs one trip round the event loop.
function makeSlowSink() {
let processed = 0;
const sink = new Writable({
objectMode: true,
highWaterMark: HIGH_WATER_MARK,
write(_chunk, _enc, done) {
processed++;
setImmediate(done);
},
});
return { sink, processed: () => processed };
}
async function withoutBackpressure() {
const { sink, processed } = makeSlowSink();
let falseCount = 0;
for (let i = 0; i < ITEMS; i++) {
// Wrong: the return value is ignored, so every item piles up in the writable buffer.
if (!sink.write(Buffer.alloc(PAYLOAD_BYTES))) falseCount++;
}
console.log("WITHOUT backpressure");
console.log(" write() returned false", falseCount, "times, and the loop ignored it");
console.log(" writableLength right after the loop:", sink.writableLength, "items");
console.log(" processed by consumer so far:", processed());
sink.end();
await once(sink, "finish");
}
async function withBackpressure() {
const { sink, processed } = makeSlowSink();
let maxQueued = 0;
let drains = 0;
for (let i = 0; i < ITEMS; i++) {
const ok = sink.write(Buffer.alloc(PAYLOAD_BYTES));
maxQueued = Math.max(maxQueued, sink.writableLength);
// Right: stop producing until the stream says it has room again.
if (!ok) {
await once(sink, "drain");
drains++;
}
}
console.log("WITH backpressure");
console.log(" highWaterMark:", HIGH_WATER_MARK, "items");
console.log(" max writableLength seen:", maxQueued, "items");
console.log(" drain events awaited:", drains);
sink.end();
await once(sink, "finish");
console.log(" processed by consumer:", processed());
}
// Load shedding: a bounded queue that refuses work instead of growing.
function sheddingDemo() {
const CAPACITY = 100;
const ARRIVALS = 1000;
const queue: number[] = [];
let shed = 0;
for (let i = 0; i < ARRIVALS; i++) {
if (queue.length >= CAPACITY) shed++; // would answer 429 or 503 + Retry-After
else queue.push(i);
}
console.log("LOAD SHEDDING");
console.log(" arrivals:", ARRIVALS, "capacity:", CAPACITY);
console.log(" accepted:", queue.length, "shed:", shed);
}
console.log("node", process.version, "items:", ITEMS, "payload:", PAYLOAD_BYTES, "bytes");
await withoutBackpressure();
await withBackpressure();
sheddingDemo();
Here is the output from running that file on Node v26.10.0. I removed memory readings from the script because process memory figures vary from run to run, so only queue counts are shown. The counts are deterministic for this script.
$ node backpressure-demo.ts
node v26.10.0 items: 100000 payload: 1024 bytes
WITHOUT backpressure
write() returned false 99985 times, and the loop ignored it
writableLength right after the loop: 100000 items
processed by consumer so far: 1
WITH backpressure
highWaterMark: 16 items
max writableLength seen: 16 items
drain events awaited: 6250
processed by consumer: 100000
LOAD SHEDDING
arrivals: 1000 capacity: 100
accepted: 100 shed: 900
Read the two blocks against each other. Ignoring the return value left all 100,000 items buffered in the stream, with the consumer having processed just one. Obeying it never let the buffer exceed 16 items, and the 6,250 drain events are simply 100,000 items divided by the 16 that fit each time. The shedding block is plain arithmetic: 1,000 arrivals into a 100-slot queue accepts 100 and refuses 900. If you use stream.pipeline or pipe, the library does this waiting for you; you only write the loop by hand when you call write yourself.
How do pull-based workers, TCP windows and reactive streams compare?
They all solve the same problem by moving who decides the rate. TCP is the original example: the receiver advertises how many octets it is willing to accept, and the sender may not exceed that, as described in RFC 9293. Reactive Streams generalises the idea by making the subscriber request a number of items, so the publisher can never send more than was asked for.
Mechanism
Who sets the rate
Signal
If the signal is ignored
Node.js write() and drain
The writable stream
write() returns false, then drain fires
Items pile up in the internal buffer
TCP receive window
The receiver
Window size advertised in each segment
Sender overruns the receiver
Load shedding
The server edge
429 or 503 with Retry-After
Client retries immediately and deepens the overload
Pull-based worker
The consumer
Consumer asks for the next batch
Not possible: nothing is sent until it asks
Reactive Streams request(n)
The subscriber
Subscriber requests a count of items
Spec violation: the publisher must not send more than requested
Pull is the simplest to reason about because the consumer controls its own intake. A queue-backed worker that claims a small batch, finishes it and only then claims more can never be flooded, however fast jobs are enqueued.
// Pull-based worker: it only asks for work when it has capacity, so the
// producer side can never push more than the workers can hold.
const BATCH = 10;
const IDLE_MS = 500;
async function runWorker(claim: (n: number) => Promise<Job[]>, handle: (j: Job) => Promise<void>) {
while (true) {
const jobs = await claim(BATCH); // e.g. SELECT ... FOR UPDATE SKIP LOCKED LIMIT 10
if (jobs.length === 0) {
await new Promise((r) => setTimeout(r, IDLE_MS)); // empty queue: back off, do not spin
continue;
}
await Promise.all(jobs.map(handle)); // finish the batch BEFORE asking for more
}
}
interface Job { id: number; payload: unknown }
The cost of pull is latency at idle. A worker that polls every 500 ms adds up to that delay to a job arriving on an empty queue, which is why long-polling or a notification channel is often paired with it.
How do you autoscale on queue depth?
Queue depth is a better scaling signal than CPU for workers, because a worker waiting on a slow downstream shows low CPU while the backlog grows. Divide arrival rate by per-worker service rate to get the workers needed, and divide the backlog by the surplus capacity to get the drain time. The figures below are assumed, to show the method.
Worker capacity: mu = 50 messages/s per worker (20 ms of work each)
Arrival rate: lambda = 200 messages/s
workers needed to hold the line = ceil(lambda / mu) = ceil(200 / 50) = 4
Backlog already queued: B = 180,000 messages, scale to 6 workers (300/s)
net drain rate = 300 - 200 = 100 messages/s
time to clear = 180,000 / 100 = 1,800 s = 30 minutes
Estimated wait for a new message = depth / (workers x mu)
at 180,000 deep with 4 workers: 180,000 / 200 = 900 s
Kubernetes lets the Horizontal Pod Autoscaler scale on external metrics, so a queue length exported by your metrics pipeline can drive replicas, provided a metrics adapter serves the external metrics API. On a single VPS the same arithmetic still helps: it tells you how many worker containers to run and whether the box can drain a backlog at all. Prefer the age of the oldest message when you can get it, because depth alone says nothing about how long users are waiting.
Which backpressure strategy should you pick?
Work through these in order. Most systems end up with more than one, because the edge, the queue and the worker each need their own protection.
Put a bound on every queue and buffer, and write down what happens when it is full: block, reject or drop.
At the public edge, cap concurrency and answer 503 with Retry-After, or 429 per client, instead of queueing past your latency budget.
Between internal workers, prefer pull-based consumption with a small batch size so the consumer sets its own pace.
In Node.js, check the return value of write() and wait for drain, or use pipeline so the library does it.
Alert on queue depth and message age, and scale workers from the arrival rate divided by per-worker throughput.
A fast producer crashes a slow consumer only when nothing in the system can say no. Bound the queue, decide what full means, and let the consumer set the pace through pull, drain or a window. Do it at every hop, because a signal that stops at one boundary just moves the pile to the next.