Backpressure in Node.js Streams, APIs, and Queues
Learn how to control producer-consumer imbalance in Node.js using stream backpressure, bounded concurrency, queue limits, admission control, retries, and load shedding.
28 min read
What Happens When Producers Are Faster Than Consumers?
Imagine a system receiving:
15,000 messages/secondbut workers can process only:
10,000 messages/secondThat leaves:
15,000 incoming
-10,000 processed
-----------------
5,000 waiting every secondSo the backlog grows by:
5,000 messages/secondAfter one minute:
5,000 × 60
= 300,000 waiting messagesAfter one hour:
5,000 × 60 × 60
= 18,000,000 waiting messagesSomething eventually has to hold those 18 million messages.
It might be:
Node.js memory
Redis
Kafka
RabbitMQ
database
diskIf they are buffered inside application memory, the process may run out of memory very quickly.
Even if they are stored safely in a durable queue, users may eventually wait hours for their work.
This is the problem backpressure tries to solve.
What Is Backpressure?
Backpressure means:
A slower consumer communicates that the producer should stop, slow down, or send less work.
The system may respond by:
pausing
waiting
limiting concurrency
queueing within a bound
rejecting requests
dropping low-value workThe important part is that the producer cannot continue producing unlimited work when the consumer cannot keep up.
Producer and Consumer Mental Model
Think of a pipeline:
Producer
↓
Buffer
↓
ConsumerFor example:
HTTP upload
↓
Node.js stream
↓
S3/object storageor:
API
↓
job queue
↓
worker
↓
PostgreSQLor:
100,000 array items
↓
Promise processing
↓
external APIThe producer creates work.
The consumer processes work.
If:
producer rate <= consumer ratethe system is usually stable.
If:
producer rate > consumer ratefor long enough, the backlog continuously grows.
Backpressure Is Not Just About Streams
Developers often hear backpressure in the context of Node.js streams.
But the same idea appears throughout backend systems.
Examples:
Node.js streams
Promise concurrency
HTTP APIs
database connections
job queues
Kafka consumers
third-party APIs
file uploads
WebSocketsThe implementation changes.
The core problem remains:
Too much work arrives
faster than downstream capacity.Why an Unlimited Buffer Does Not Solve the Problem
A common response is:
"Just queue everything."That only moves the problem.
Imagine:
arrival rate = 500 jobs/minute
worker capacity = 50 jobs/minuteEvery minute:
450 additional jobsare added to the queue.
After one hour:
27,000 jobs waitingThe queue prevents immediate memory exhaustion.
But the system is still overloaded.
Eventually:
job latency becomes huge
queue storage grows
retries accumulate
users stop caring about old work
recovery takes hoursAn unbounded queue converts:
immediate overloadinto:
future overload + enormous latencyBackpressure requires a finite limit somewhere.
Node.js Streams Already Support Backpressure
Node.js streams have built-in backpressure behavior.
Suppose we are writing chunks into a writable stream:
for (const chunk of source) {
destination.write(chunk);
}This looks simple.
But what happens if:
source produces data very quicklywhile:
destination writes slowlyNode starts buffering data.
If we keep writing indefinitely, memory usage grows.
write() Gives You a Backpressure Signal
A writable stream's:
destination.write(chunk);returns a boolean.
If it returns:
true;the writable stream still has capacity.
If it returns:
false;its internal buffer has reached its high-water mark.
That means:
Stop writing for now.
Correct Stream Handling
for (const chunk of source) {
if (!destination.write(chunk)) {
await once(destination, "drain");
}
}
destination.end();The important part is:
if (!destination.write(chunk))When it returns false, we wait for:
drainbefore continuing.
What Does drain Mean?
Imagine the writable stream's buffer looks like:
Buffer
████████████████████
FULLwrite() returns:
falseThe downstream destination continues processing.
Eventually the buffer becomes:
████████Node emits:
drainThat tells the producer:
You can start writing again.So the flow becomes:
write
write
write
buffer fills
↓
pause
↓
consumer drains buffer
↓
"drain"
↓
resumeThat is backpressure.
What Happens If You Ignore write() === false?
Bad:
for (const chunk of source) {
destination.write(chunk);
}If the producer is much faster than the consumer:
Producer
██████████████████████████
Consumer
██████the difference accumulates in memory.
Memory may look like:
100 MB
300 MB
800 MB
1.5 GB
...Eventually:
process slows
garbage collection increases
latency rises
process may crashThe stream did tell you it was overloaded.
You ignored the signal.
highWaterMark Is Not a Hard Memory Limit
Node streams use a setting called:
highWaterMarkThis tells the stream roughly when it should start applying backpressure.
It does not mean:
Node can never buffer more than this value.Think of it more as:
"At this point, producers should stop pushing."If you ignore the backpressure signal, memory can still continue growing.
Prefer pipeline()
For common stream composition, Node provides:
pipeline();Example:
await pipeline(readable, transform, writable);Conceptually:
Readable
↓
Transform
↓
Writablepipeline() helps coordinate:
backpressure
errors
completion
cleanupThat is usually safer than manually wiring many:
pipe()
error
close
finishlisteners yourself.
Example: Large File Processing
Suppose we process a 10 GB file.
Bad design:
Read entire file
↓
store in memory
↓
process
↓
uploadObviously this requires huge memory.
Streams instead allow:
Read chunk
↓
transform chunk
↓
write chunkwhile keeping memory bounded.
Backpressure ensures:
disk read speeddoes not greatly outrun:
destination write speedPromise Concurrency Has the Same Problem
Backpressure problems are not limited to streams.
Consider:
await Promise.all(items.map(processItem));This is common.
It may work perfectly for:
10 itemsor:
100 itemsBut what about:
100,000 items?
What Promise.all() Actually Does
Suppose:
items.map(processItem);contains 100,000 items.
If processItem() immediately starts async work, you may create:
100,000 promises
100,000 DB requests
100,000 HTTP requests
100,000 pending closuresalmost immediately.
The code is effectively saying:
Start everything now.
That can overwhelm:
database connection pool
HTTP socket pool
provider quota
CPU
memory
RedisExample: Database Pool
Imagine PostgreSQL pool size:
20 connectionsThen your code creates:
100,000 queriesusing:
Promise.all(...)Only about 20 can run concurrently.
The remaining work waits somewhere:
application memory
pool wait queue
driver internalsYou have not increased database throughput.
You have only created a huge waiting list.
Bound Promise Concurrency
Instead of starting every item immediately, use a fixed number of workers.
Example:
const workers = Array.from({ length: 20 }, async () => {
while (true) {
const item = nextItem();
if (!item) {
return;
}
await processItem(item);
}
});
await Promise.all(workers);Now maximum active work is approximately:
20 operationsrather than:
100,000 operationsWorker Pool Mental Model
Without bounds:
100,000 items
↓
100,000 concurrent operations
↓
databaseWith a worker pool:
100,000 items
↓
20 workers
↓
databaseThe input list can still be large.
But active work is controlled.
How Do You Choose Concurrency?
Do not automatically choose:
20because an example used 20.
Concurrency should come from downstream capacity.
Suppose:
database pool = 30 connectionsbut your application also uses those connections for:
authentication
normal API traffic
background jobsGiving one bulk operation all 30 connections may starve the rest of the system.
You might choose:
bulk worker concurrency = 10and leave capacity for other operations.
Latency Also Matters
Suppose an external API supports:
100 requests/secondand each request takes approximately:
200 msA concurrency of around:
20can theoretically sustain approximately:
20 / 0.2
= 100 requests/secondbefore other overhead.
A useful relationship is:
throughput ≈ concurrency / latencyThis is not exact, but it helps reason about capacity.
Then measure real behavior.
More Concurrency Does Not Always Mean More Throughput
Imagine increasing workers:
10
20
50
100
500At first throughput may rise.
Eventually the database becomes saturated.
Then increasing concurrency causes:
more waiting
higher p95 latency
higher p99 latency
more timeouts
more memory
more retrieswhile throughput barely changes.
You may see:
20 workers → 1,000 jobs/sec
50 workers → 1,400 jobs/sec
100 workers → 1,450 jobs/sec
500 workers → 1,300 jobs/secMore concurrency actually made the system worse.
API Admission Control
Now imagine an HTTP endpoint:
POST /reportsEach report generation is expensive.
Suppose the backend can safely run:
50 reports concurrentlyThen 5,000 requests arrive.
If we simply accept all requests:
50 running
4,950 waitingYou did not increase throughput.
You created a huge wait queue.
What Is Admission Control?
Admission control asks:
Should this work even enter the system right now?
Instead of always accepting requests:
incoming request
↓
capacity available?
┌───┴───┐
yes no
↓ ↓
accept rejectThis protects expensive downstream resources.
Bounded Semaphore
A concurrency gate may look conceptually like:
incoming
↓
concurrency gate
↓
workers
↓
database/providerSuppose:
active limit = 50
waiting queue = 100If:
50 active
+
100 waitingalready exist, request 151 gets rejected.
Now system memory and waiting time remain bounded.
Why Bound the Waiting Queue Too?
Suppose concurrency is bounded to:
50but waiting queue is unlimited.
You still have:
50 running
100,000 waitingThe running work is bounded.
But memory and latency are not.
So both should usually be finite:
active operations
+
waiting operationsWhen Should the API Return 429?
Use:
429 Too Many Requestswhen the client exceeded a defined policy.
For example:
tenant may run only
2 reports at onceor:
tenant exceeded
100 requests/minuteThat is a client or account-level limit.
When Should the API Return 503?
Use:
503 Service Unavailablewhen the service itself temporarily lacks capacity.
For example:
all report workers saturated
database pool overloaded
queue full globally
dependency unavailableEasy distinction:
429
→ you exceeded your allowed share503
→ the service currently cannot accept more workThere can be product-specific exceptions, but this is a useful default mental model.
Rejecting Early Can Be Better Than Waiting
Imagine estimated wait time is:
30 secondsbut client timeout is:
5 secondsAccepting the request makes no sense.
The client will leave before work starts.
Better:
503 Service Unavailable
Retry-After: 10Now the request fails quickly without consuming resources for useless waiting.
Queue Backpressure
Durable queues are useful because they move waiting work out of Node.js memory.
For example:
API
↓
Kafka / RabbitMQ / BullMQ
↓
Workers
↓
DatabaseIf workers temporarily slow down, jobs remain durable.
This is much safer than storing thousands of pending promises in application memory.
But queues are not infinite.
Durable Does Not Mean Unlimited
Suppose:
producer = 500 jobs/minute
consumer = 50 jobs/minuteThe durable queue survives.
But backlog still grows:
+450 jobs/minuteEventually:
storage grows
latency becomes unacceptable
messages become stale
recovery becomes difficultSo a queue changes where backlog lives.
It does not remove the capacity problem.
Monitor Arrival Rate and Completion Rate
Two of the most important numbers are:
publish rate
completion rateFor example:
Publish:
500 jobs/minute
Completion:
480 jobs/minuteBacklog grows by:
20 jobs/minuteMaybe that is acceptable for a short burst.
But if it continues for hours, it is not.
Queue Depth Is Not Enough
Suppose Queue A has:
10,000 jobsand each takes:
1 msEstimated drain time is roughly:
10 secondsMaybe healthy.
Now Queue B has:
100 jobsbut each job waits:
1 hourThat may be a major incident.
So monitor:
queue depth
+
oldest job agenot queue depth alone.
Consumer Lag
In systems like Kafka, a related metric is:
consumer lagConceptually:
latest produced offset
-
latest processed offset
=
consumer lagGrowing lag tells us:
consumers are falling behind producersBut again, time is important.
A lag of:
100,000 messagesmay be tiny for a system processing:
1,000,000 messages/secwhile a lag of:
5,000may be serious for a system processing:
10 messages/secTrack lag in both:
messages
and
timewhere possible.
Can We Just Add More Consumers?
Sometimes.
Suppose:
10 workerscannot keep up.
You increase to:
20 workersand throughput doubles.
Great.
But eventually something else becomes the bottleneck.
For example:
workers
↓
PostgreSQLIf PostgreSQL supports:
200 concurrent operations safelystarting:
1,000 workersdoes not solve the problem.
The database becomes overloaded.
Scaling Consumers Moves Pressure Downstream
Think of the pipeline:
Queue
↓
Workers
↓
DatabaseIf you increase workers aggressively:
Queue pressure ↓
Database pressure ↑The bottleneck moves.
You must consider the entire chain.
Partition Count Can Also Limit Throughput
For Kafka-style systems, suppose:
topic partitions = 10and one consumer in a consumer group actively handles each partition.
Adding:
100 consumersdoes not necessarily give:
10× more throughputbecause only roughly:
10 consumersmay have active partitions.
Scaling requires understanding:
consumer count
partition count
downstream capacitytogether.
Retry Traffic Is Also Traffic
Imagine normal traffic is:
10,000 jobs/secThe system handles it.
Then a dependency fails.
5,000 jobs fail and retry immediately.
Now effective input becomes:
Original traffic:
10,000/sec
Retry traffic:
5,000/sec
Total:
15,000/secThe outage already reduced capacity.
Retries just increased demand.
This can turn a small dependency problem into a full system outage.
Retry Storm Example
Imagine:
Producer:
10,000/sec
Consumer during outage:
6,000/secDeficit:
4,000/secNow failed jobs retry instantly.
Effective producer rate may become:
14,000/secor higher.
Deficit grows again.
This feedback loop is sometimes called a:
retry stormUse Backoff
Instead of:
fail
↓
retry immediately
↓
fail
↓
retry immediatelyuse increasing delays.
For example:
1 second
2 seconds
4 seconds
8 seconds
16 secondsThis gives the downstream system time to recover.
Add Jitter
If 100,000 jobs all fail at the same moment and use:
retry exactly after 10 secondsthen after 10 seconds:
100,000 retries arrive togetherThat creates another spike.
Add randomness:
8.2 sec
9.4 sec
10.1 sec
11.7 sec
...This spreads retries over time.
That is jitter.
Bound Retries Too
Do not retry forever.
For example:
attempt 1
attempt 2
attempt 3
attempt 4
↓
dead-letter queueor:
mark failedOtherwise broken jobs can become permanent producers of traffic.
Decide What Can Be Dropped
Backpressure is not only a technical question.
It is also a product decision.
Suppose the system is overloaded.
Different work has different value.
For example:
Never drop:
confirmed payment record
Delay:
invoice email
Sample:
low-value analytics
Reject:
optional reportThis priority should be decided intentionally.
Critical Work vs Optional Work
Imagine resources are nearly exhausted.
Would you rather preserve:
POST /payments/confirmor:
POST /analytics/exportProbably the payment path.
So the system might do:
Critical:
continue accepting
Normal:
slow down
Optional:
rejectThis is load shedding.
What Is Load Shedding?
Load shedding means deliberately rejecting work to keep the rest of the system healthy.
It sounds bad because:
we are rejecting usersBut compare:
Option A
Accept everything:
10,000 requests accepted
↓
everything slows
↓
database collapses
↓
almost every request failsOption B
Reject 30% immediately:
7,000 requests succeed normally
3,000 receive fast retryable errorsOption B may provide much better overall reliability.
Overload Should Be Explicit
A common bad design is:
keep accepting
until memory,
connections,
or database
finally collapseThat is accidental overload behavior.
A better system decides:
At capacity X,
we start rejecting Y.That is deliberate load shedding.
End-to-End Backpressure
Backpressure has to exist through the entire pipeline.
Suppose:
HTTP upload
↓
Node.js
↓
object storageIf object storage slows down, Node should stop reading from the client as quickly.
Otherwise:
client uploads fast
↓
Node buffers huge amount
↓
storage writes slowly
↓
memory growsA properly streamed pipeline propagates the pressure backward.
Database Example
Suppose queue consumers write to PostgreSQL.
Kafka
↓
consumer
↓
PostgreSQLDatabase pool:
30 connectionsIf you process:
500 messages concurrentlyyou create:
30 running DB operations
470 waiting operationsInstead, consumer concurrency should be related to pool capacity.
For example:
consumer DB concurrency = 15leaving room for:
API traffic
admin operations
other workersThird-Party API Example
Suppose a provider allows:
50 requests/secondbut your queue contains:
100,000 jobsRunning 1,000 workers only causes:
429 responses
timeouts
retries
more queue trafficThe correct backpressure may be:
queue
↓
rate limiter
↓
50 requests/sec
↓
providerQueue size and provider limit need to work together.
Backpressure vs Rate Limiting
These concepts overlap but are different.
Rate limiting asks:
How much work is someone allowed to create over time?
Example:
100 requests/minuteBackpressure asks:
What happens when downstream capacity cannot currently keep up?
Example:
database saturated
→ stop accepting more expensive workYou may use both.
Backpressure vs Concurrency Limiting
Concurrency limiting asks:
How much work can run at once?
Example:
20 database operations concurrentlyThis is one common mechanism used to implement backpressure.
Backpressure is the broader concept.
Backpressure vs Queueing
Queueing says:
work can waitBackpressure says:
waiting must remain controlledA queue can be part of backpressure.
An unlimited queue is usually not a complete backpressure strategy.
Backpressure vs Load Shedding
Backpressure tries to:
slow or regulate incoming workLoad shedding says:
capacity is exhausted
→ reject some workLoad shedding is often the final step when backpressure cannot reduce demand enough.
A Healthy Bounded Pipeline
A production pipeline might look like:
Incoming traffic
↓
rate limit
↓
admission control
↓
bounded queue
↓
bounded worker concurrency
↓
databaseAt every step there is a finite capacity.
For example:
API rate:
1,000 requests/sec
Concurrent requests:
200
Queue:
maximum 5,000
Workers:
50
DB pool:
100The exact numbers depend on the system.
The important idea is:
nothing is infinitely bufferedWhy Unbounded Systems Fail Slowly
One dangerous property of unbounded queues is that everything may initially look fine.
Suppose:
arrival = 110/sec
capacity = 100/secOnly:
10 jobs/secare added to the backlog.
After one minute:
600 jobsMaybe nobody notices.
After one hour:
36,000 jobsNow latency is large.
After six hours:
216,000 jobsEventually the system appears to fail "suddenly."
But overload started hours earlier.
This is why monitoring the difference between arrival and completion rates is useful.
Queue Age Can Be More Important Than Queue Size
Imagine an email queue.
Normal requirement:
email delivered within 30 secondsCurrent state:
queue depth = 2,000That might be fine if workers clear it in five seconds.
But if:
oldest email = 45 minutes oldthe service is unhealthy.
The product cares about:
how long users waitnot only:
number of queued jobsCapacity Should Be Measured in Useful Work
Do not only ask:
How many requests can Node handle?Ask:
How many successful business operations
can the full system complete per second?For example:
Node:
20,000 req/sec
PostgreSQL:
3,000 writes/sec
Payment provider:
500 charges/secFor payment creation, actual system capacity may be closer to:
500/secbecause that is the tightest downstream dependency.
The Slowest Dependency Often Defines Throughput
Imagine:
API → 10,000/sec
Redis → 20,000/sec
Database → 2,000/sec
Provider → 500/secIf every request requires the provider:
effective maximum
≈ 500/secIncreasing API worker threads does not fix that.
Backpressure should protect the provider boundary.
Deadlines Matter Too
Suppose the queue wait is:
20 secondsbut the user only cares about the result for:
5 secondsThe queued work may already be useless.
Use:
request deadlines
job deadlines
queue-age limitsto avoid performing work after its value has expired.
Queue Age Limits
For example:
report preview jobs
maximum useful age = 30 secondsA worker picks up a job after:
90 secondsInstead of spending CPU on it:
discard or fail itbecause the result is no longer useful.
This prevents old work from delaying new useful work.
Async Work Can Use 202 Accepted
If work may wait for minutes, do not keep an HTTP request open.
For example:
POST /reportscan return:
202 Acceptedwith:
{
"operationId": "op-123",
"status": "pending"
}Then:
API
↓
bounded queue
↓
workerhandles the expensive work asynchronously.
This makes queueing explicit to the user.
But Even Async Queues Need Admission Control
Suppose the job queue already contains:
1,000,000 report requestsAccepting another million because:
"It's asynchronous."does not solve the problem.
The API may need to say:
queue currently at capacityand reject or delay new optional work.
Autoscaling Does Not Remove Backpressure
Suppose Kubernetes sees:
high queue depthand scales workers:
10
→ 100That can help if downstream capacity exists.
But if PostgreSQL is already at maximum:
10 workers
→ DB 80% load
100 workers
→ DB 100% load
→ timeoutsAutoscaling can make the incident worse.
Scale based on the whole system, not only queue size.
Backpressure in WebSockets
Imagine clients send:
1,000 messages/secbut the server can send to a slow device at:
10 messages/secIf you keep buffering outbound messages per socket:
buffer grows foreverA production design may need:
bounded per-client buffer
drop stale updates
disconnect extremely slow clients
coalesce repeated state updatesFor example, if 500 stock-price updates are queued for the same symbol, maybe only the newest one matters.
Backpressure in Event Systems
Suppose telemetry events arrive faster than they can be stored.
Different event classes may have different policies:
Payment event
→ never drop
Security event
→ preserve
Page-view analytics
→ sample
Debug telemetry
→ drop firstBackpressure works better when the system understands work value.
Production Metrics
A backpressure system should be observable.
Useful metrics include:
incoming rate
completion rate
queue depth
oldest job age
consumer lag
active concurrency
waiting concurrency
rejection rate
429 rate
503 rate
retry rate
dead-letter rate
database pool usage
provider latency
memory usageThe Most Important Relationship
Watch:
arrival rate
vs
completion rateIf:
arrival > completionfor a sustained period:
backlog growsIf:
completion > arrivalthe system can drain backlog.
Example
Arrival:
1,200 jobs/sec
Completion:
1,000 jobs/secBacklog growth:
+200/secThen after scaling:
Arrival:
1,200/sec
Completion:
1,500/secBacklog drains at:
300/secIf backlog is:
30,000 jobsapproximate drain time:
30,000 / 300
= 100 secondsThis kind of calculation is useful during incidents.
Test Slow Consumers
A common testing mistake is testing only:
fast database
fast Redis
fast external provider
fast diskEverything works.
Production problems occur when a consumer slows down.
Test situations like:
database latency × 10
provider starts returning 429
object storage becomes slow
worker throughput drops by 50%
network stalls
queue fillsThen verify that:
memory stays bounded
concurrency stays bounded
requests reject predictably
critical work continuesTest Burst Traffic
Average traffic is not enough.
Suppose average is:
1,000 req/secbut every minute a batch client creates:
20,000 requests
in two secondsAverage capacity may look fine.
Burst capacity may not be.
Load tests should include:
steady traffic
short bursts
long bursts
slow consumers
dependency outages
recoveryTest Recovery Too
Suppose the database is slow for five minutes.
The queue grows.
Then database speed returns to normal.
Question:
Can the system recover?Or do retries and backlog keep it overloaded?
Sometimes after the dependency recovers:
huge backlog
+
huge retry waveimmediately overloads it again.
Recovery needs controlled drain rates.
Avoid the Thundering Herd During Recovery
Imagine 100,000 jobs are waiting.
Database recovers.
If every worker immediately runs at full speed:
100,000 jobs
↓
databaseit may fail again.
Instead ramp up:
consumer concurrency
10
→ 20
→ 40or use rate-limited draining.
Recovery is part of backpressure design.
A Practical API + Queue Example
Suppose we build a video-report system.
Incoming traffic:
POST /reportsEach job takes:
10 secondsWorker capacity:
20 jobs concurrentlyA reasonable design might be:
POST /reports
↓
tenant rate limit
↓
global admission control
↓
queue capacity check
↓
persist job
↓
202 AcceptedThen:
queue
↓
20 workers
↓
database + external providersIf queue is full:
503 Service Unavailableor perhaps:
429if the tenant exceeded its own allowed queue share.
This is more predictable than accepting unlimited jobs.
A Practical Stream Example
Suppose an API receives a large upload and forwards it to object storage.
Healthy flow:
Client upload
↓
Readable request stream
↓
Transform
↓
Object-storage writableWhen storage slows:
writable buffer fills
↓
write() returns false
↓
upstream reading slowsPressure propagates backward toward the client.
This keeps application memory bounded.
A Practical Database Import Example
Suppose you import:
1,000,000 CSV rowsBad:
await Promise.all(rows.map(insertRow));Better:
CSV stream
↓
parser
↓
bounded batch
↓
worker pool
↓
PostgreSQLFor example:
batch size = 500
worker concurrency = 5Then adjust based on real database throughput.
Common Backpressure Mistakes
Mistake 1: Using Promise.all() Everywhere
Fine:
5 independent operationsDangerous:
500,000 operationsBound large workloads.
Mistake 2: Ignoring write() === false
Node is explicitly telling you:
slow downRespect it.
Mistake 3: Assuming a Queue Solves Capacity
A queue stores backlog.
It does not create processing capacity.
Mistake 4: Unlimited Waiting Queues
Bounded concurrency with an unlimited waiting list can still exhaust memory and create massive latency.
Bound both.
Mistake 5: Scaling Workers Without Checking the Database
You may simply move overload from:
queueto:
PostgreSQLMistake 6: Immediate Retries
Retries create new load.
Use:
bounded attempts
backoff
jitterMistake 7: No Work Priority
If optional exports consume all capacity, critical payment work may fail.
Classify work.
Mistake 8: Looking Only at Queue Size
Track:
oldest work agetoo.
Mistake 9: Waiting Longer Than the Work Is Useful
If the client deadline is gone, continuing may waste capacity.
Use deadlines and queue-age policies.
A Simple Mental Model
When building a pipeline, ask:
1. Who produces the work?Then:
2. Who consumes it?Then:
3. What is the maximum sustainable
consumer rate?Then:
4. Where does excess work wait?Then:
5. How much waiting is allowed?Then:
6. What happens when that limit
is reached?Then:
7. Which work is critical,
delayed, sampled, or rejected?Then:
8. How will retries affect
the input rate?Finally:
9. How will we know the system
is falling behind?If you cannot answer:
Where does excess work go?you probably have an unbounded buffer somewhere.
Backpressure Across a Whole System
A mature system may look like:
Clients
↓
Edge rate limit
↓
API admission control
↓
bounded request concurrency
↓
bounded durable queue
↓
bounded workers
↓
provider/database limiter
↓
dependencyEach boundary protects the next one.
For example:
Edge
→ protects API
API admission control
→ protects queue
Queue limits
→ protect storage + latency
Worker concurrency
→ protects database
Provider rate limiter
→ protects external quotaThat is end-to-end flow control.
Production Checklist
Before shipping a high-throughput Node.js workflow:
□ Identify producers and consumers.
□ Measure sustainable consumer rate.
□ Respect stream write() backpressure.
□ Wait for drain when required.
□ Prefer pipeline() for stream chains.
□ Do not create huge Promise.all batches.
□ Bound async concurrency.
□ Match worker concurrency to
downstream capacity.
□ Bound waiting queues.
□ Add admission control for
expensive synchronous work.
□ Use 429 for client-specific limits.
□ Use 503 for temporary service
capacity exhaustion where appropriate.
□ Keep durable queues finite.
□ Monitor publish and completion rate.
□ Monitor queue depth.
□ Monitor oldest-work age.
□ Monitor consumer lag.
□ Include retry traffic
in capacity calculations.
□ Use bounded retry attempts.
□ Add exponential backoff.
□ Add jitter.
□ Define work priorities.
□ Decide what may be delayed.
□ Decide what may be sampled.
□ Decide what may be rejected.
□ Preserve critical work.
□ Add queue-age/deadline rules.
□ Do not autoscale beyond
downstream capacity.
□ Test slow databases.
□ Test slow external providers.
□ Test burst traffic.
□ Test queue saturation.
□ Test recovery after outages.
□ Keep memory and concurrency bounded.Conclusion
Backpressure exists because producers and consumers rarely run at exactly the same speed.
If:
producer > consumerfor long enough, something must grow:
memory
queue
latency
disk usage
retry trafficWithout a limit, that growth eventually becomes a failure.
In Node.js streams:
write() === false
→ slow downWith asynchronous work:
huge Promise.all
→ replace with bounded concurrencyFor APIs:
capacity full
→ stop admitting unlimited workFor queues:
consumer lag grows
→ throttle producers,
scale carefully,
or reject lower-value workFor retries:
failure
→ backoff + jitterAnd during severe overload:
protect critical work
→ shed lower-priority workThe most important idea is:
Backpressure is not about making the producer slower all the time. It is about preventing the producer from overwhelming the consumer when their capacities no longer match.
A healthy production system therefore does not rely on infinite buffers.
It uses:
bounded memory
bounded concurrency
bounded queues
admission control
rate limits
deadlines
retries with backoff
priority
load shedding
observabilityto keep overload controlled.
The final question to ask at every system boundary is:
If the next component becomes slower than this one, where does the extra work go, and what stops that backlog from growing forever?
If the answer is clear and bounded, you have a real backpressure strategy.
For the lower-level stream primitives behind write(), drain, pipeline(),
and highWaterMark, start with the
Node.js streams guide. For overload at the
whole-API level, continue with what breaks between 100 and 10,000 requests per
second.
References
Related
Written by
Faisal
Software engineer writing about backend systems, Node.js, system design, scalable applications, and modern web and mobile development.