Hiring BarSupport

Backpressure & Overload Control

Pattern45 min read5 diagrams

Technologies that implement this pattern: Apache Kafka · Apache Flink · API Gateway · Kubernetes · Redis

Why This Matters#

Every outage that starts with "the queue grew" is a backpressure failure. Somewhere a producer was allowed to send work faster than a consumer could finish it, the difference went into a buffer, and the buffer was the only part of the system nobody had sized. Queues do not create capacity. They move the overload forward in time and convert it into latency, memory, and work that will be thrown away because the caller stopped waiting minutes ago.

Most candidates treat overload as a capacity question: "autoscale the consumers." Staff engineers treat it as a signal-propagation question. When the slowest stage is full, who finds out, how fast, and what do they do about it? There are only four answers — block the producer, slow the producer, buffer the excess, or drop it — and every system picks one at every hop whether or not anyone decided. The design question is not how big is the queue. It is where does "I'm full" travel, and what does the work at the edge do when it arrives.

The second reframe: backpressure and load shedding are different tools with different victims. Backpressure slows the producer, so the cost lands on whoever is upstream: a batch job that runs longer, a client that waits, a Kafka topic that grows. Shedding drops the work, so the cost lands on a user who gets an error or a degraded answer. A design that cannot say which one it is doing at each hop, and who signed off on that, will do both at random during the incident.

If you can walk an interviewer from "where is the bottleneck" to "how does it signal upstream" to "where does the signal stop and turn into a drop, and who owns that decision," you are answering at Staff level.

The 60-Second Version#

  • Every queue is bounded — the only question is whether you chose the bound. An unbounded in-memory queue is bounded by the heap; it fails with an OOM at the worst possible moment. Size queues by time, not count: max_queue_len = drain_rate × max_acceptable_wait. At 2,000 req/s and a 100ms wait budget, that's 200 slots, not 10,000.
  • Little's law sets the concurrency you can afford. in_flight = throughput × latency. A service doing 2,000 req/s at 50ms needs ~100 in flight. If in-flight is 1,000, nine-tenths of it is queueing — latency, not work.
  • Pull beats push when the consumer is the bottleneck. In a pull (or credit-based) protocol the consumer asks for n items and the producer can never outrun it. Push needs an explicit stop signal, and that signal is usually the part that's missing.
  • Backpressure must reach a place that can absorb it. A durable log (Kafka, with days of retention) can absorb hours of slowdown. A user's HTTP request can absorb ~1–2 seconds. If the signal arrives at a user-facing edge, it must become a fast rejection (429/503), not a hang.
  • Measure lag in seconds, not messages. lag_seconds = lag_messages / consume_rate. 5M messages is 50 seconds at 100K/s or 14 hours at 100/s. Alert on time-to-drain and on the age of the oldest unprocessed item.
  • Recovery needs headroom: drain_time = backlog / (consume_rate − produce_rate). With 20% headroom, a 1-hour outage takes ~5 hours to drain. With 5%, ~20 hours. Headroom is how you pay for backpressure.

The Problem#

Producers and consumers almost never run at the same speed. A checkout spike sends 8× normal order events into the fulfilment pipeline; a slow database makes a consumer process 300 msg/s instead of 3,000; a downstream ML scorer gets a cold cache after deploy and its p99 jumps from 40ms to 900ms. In each case the stage in front of the bottleneck has three choices: wait, hold, or throw away. If nothing is designed, it holds — in a thread pool queue, a connection pool, an HTTP/2 stream, a broker partition, a retry loop — until memory runs out or every queued item is already past its deadline. Then the system is busy doing work nobody wants, the goodput drops to near zero while CPU reads 100%, and the overload outlives the spike that caused it. The job is to make "full" an explicit signal, carry it upstream to a place that can afford to wait, and convert it into a deliberate drop at the point where waiting stops being affordable.


Case Studies That Use This Pattern#

  • Message Broker — Broker as the shock absorber; consumer lag, retention as the backpressure budget, and when the producer must block
  • Stream Processing — Credit-based flow control between operators; a slow sink throttles the source instead of crashing the job
  • Rate Limiter — Static admission at the edge; the coarse, per-tenant version of "slow down"
  • Circuit Breakers — Retry budgets and adaptive throttling so callers stop amplifying overload
  • Notifications — Fan-out bursts of 10–100× into rate-limited providers (APNs, SMS) that push back hard
  • Ad Click Aggregation — Ingest must never block the click path, so the log absorbs it and aggregation lags instead
  • Realtime WebSockets — Slow clients on a push channel: per-connection send buffers and the decision to disconnect
  • Batch & Stream Pipelines — Batch and CDC pipelines where downstream throttling is the normal state, not the exception

Which Problem Are We Solving?#

"Handle overload" hides three goals that lead to different designs. Name them and commit before drawing a box.

IntentConstraintStrategyFailure ModeCorrectness Bar
Lossless async pipeline (events, CDC, billing records)Every item must eventually be processed; latency is flexible (minutes–hours)Durable log + pull consumers + lag-based autoscaling; producer blocks only if the log itself is fullLag grows faster than you can drain; retention expires before consumptionZero loss; bounded time-to-process, product-signed
Latency-bound request path (APIs, RPC fan-out)Caller waits ≤ 1–2 s; stale work is worthlessBounded queues by time, adaptive concurrency limits, fast rejection, deadline propagationQueues fill with expired requests; goodput collapses while CPU is peggedBounded p99; explicit 429/503 rather than timeouts
Fair shared infrastructure (multi-tenant brokers, platform APIs)One tenant's burst must not slow everyonePer-tenant quotas that throttle (delay) rather than error; separate queues per tenant classNoisy neighbour fills the shared queue; head-of-line blockingPer-tenant throughput guarantee

🎯 Staff Move: "I'll treat the ingest path as intent one — lossless, minutes of lag acceptable — so backpressure terminates at the Kafka topic and nothing upstream ever waits on the consumers. The synchronous API in front of it is intent two: if the topic write itself slows past 200ms, the API rejects with a 503 and a Retry-After instead of holding the connection."


The Four Responses to Overload#

At every hop, an overloaded stage does exactly one of four things to new work. The skill is choosing it per hop on purpose.

ResponseMechanismCost Lands OnRight When
Block (synchronous backpressure)Bounded queue put() waits; TCP zero window; Kafka producer max.block.msThe producer's thread and its callersProducer is a batch job or a durable log that can wait
Slow / throttle (signalled backpressure)Credits, request(n), quota delays, 429 + Retry-AfterProducer throughput, gracefullyProducer is cooperative and can pace itself
BufferQueue, log, spillover storageLatency and memory; work goes staleBurst is short relative to buffer-in-time
Drop (shedding)Reject at admission, drop oldest, drop lowest priorityA user or a data consumer, visiblyWaiting has no value: the deadline has passed or the item is superseded
Diagram: The Four Responses to Overload

The line between backpressure and shedding is the point in the call chain where waiting stops being cheaper than failing. Backpressure travels upstream until it reaches either a durable buffer (which absorbs it as lag) or a human (who cannot absorb more than a second or two). At that boundary it must turn into a drop. Who decides what gets dropped — criticality, tenant, age — is the territory of the Graceful Degradation: Fail Open or Closed; this page is about getting the signal to that boundary quickly and losslessly.


The Core Tradeoff#

StrategyWhat WorksWhat BreaksWho Pays
Unbounded bufferNever rejects; simpleOOM or hours of stale work; latency unbounded; overload outlives the spikeOn-call at 3am; users waiting on work already abandoned
Bounded queue + blockLossless; producer naturally pacedBlocking propagates to threads that serve unrelated work; deadlocks in cyclesUnrelated callers sharing the blocked thread pool
Credit / pull-based flow controlProducer cannot outrun consumer by construction; per-channel isolationNeeds protocol support end to end; one non-compliant hop breaks itPlatform team that owns the protocol and client libraries
Static concurrency or rate limitPredictable, easy to reason aboutWrong the day capacity changes (deploy, cache cold, instance type)Users rejected when there was capacity; or the backend when there wasn't
Adaptive concurrency limitTracks real capacity from latency; no hand-tuningOscillation, slow convergence, mistakes latency of a dependency for own overloadSome false rejections during probing; on-call debugging "why 503?"
Durable log as shock absorberDecouples producer from consumer entirely; hours of slackLag becomes the failure mode; retention is a hard deadline; data goes stale silentlyData consumers seeing late results; storage bill
Shed at admissionProtects goodput; fast honest failureUsers see errors; needs priority to avoid dropping the wrong thingUsers, and product who must sign off on which ones

Staff Default Position#

Bound every queue in time, make "full" an explicit signal, and terminate backpressure at a durable log or convert it to a fast rejection — never let it end in an unbounded buffer.

The default stack: synchronous request paths use an adaptive concurrency limit per service and deadline propagation, with queue length sized to ≤ 100–200ms of drain time and rejection (429/503 with Retry-After) beyond it. Asynchronous work is handed off to a durable log at the earliest point where the user no longer needs the result; consumers pull with a bounded prefetch (max.poll.records, request(n), prefetch count) so they cannot be flooded, and they scale on lag in seconds. Retries are budgeted (≤ 10% of traffic per client) so rejection is not undone by retry amplification. Every hop has a named answer to "what happens when the next hop is full," and every terminal drop has a named owner.


When to Deviate#

  • Work cannot be dropped and cannot be durably buffered — e.g., a payment authorization in flight. Block with a short deadline and fail the whole request upstream; never silently drop a half-done transaction. Idempotency keys make the retry safe (Idempotency).
  • Superseded data — location updates, metrics gauges, presence, video frames. Drop the oldest on overflow (keep-latest), not the newest. Backpressuring a GPS stream just makes every position late.
  • Very small systems — Below ~500 req/s on a single service with slack capacity, a static concurrency limit and a fixed small queue are enough. Adaptive limiters add a debugging surface you don't need yet.
  • Third-party producers you can't slow — webhooks, partner pushes, device fleets. You cannot backpressure someone else's retry loop; accept fast into a durable log and absorb it as lag, or publish a quota with a 429 contract and enforce it at the edge (Rate Limiter).

One Question, Three Levels#

BehaviorSenior (L5)Staff (L6)Principal (L7)
First move"Put a queue between them and autoscale the consumers""Where's the bottleneck, how does it signal upstream, and where does that signal turn into a drop?""Which overload signals cross team boundaries, and is there one protocol for them — or does every team invent its own 503?"
Queue sizingLarge queue so nothing is lostSize by time: drain rate × wait budget; reject beyond it; alert on oldest-item ageSets an org-wide rule: no unbounded in-memory queues in serving paths; enforced in the shared RPC/queue libraries
LimitsStatic RPS limit per endpointAdaptive concurrency limit from latency + retry budgets so rejection isn't undoneFunds the limiter as a platform default; prices false-rejection rate against outage minutes avoided
Async lagAlert when lag > 1M messagesAlert on lag-seconds and time-to-drain; plan consumer headroom so a 1h outage drains in < 4hTreats retention and drain headroom as a budget line; decides which pipelines get guaranteed catch-up capacity
Failure"The consumer was slow, we added instances""The queue held 40 minutes of expired requests; we were 100% busy at 0% goodput"Designs the org's overload posture: which tiers block, which shed, and game-days the boundary
OwnershipWhoever owns the consumerProducer owns pacing contract; consumer owns lag SLO; on-call owns drop decisionsWrites the cross-team contract: producers must honour Retry-After and credits, consumers must publish capacity
Why "First move" separates levels

The L5 answer — queue plus autoscaling — is often right for the steady state and still gets downleveled, because autoscaling takes 1–5 minutes and the burst is over in 30 seconds, or the bottleneck is a database that doesn't autoscale at all. The Staff candidate first finds the bottleneck and traces where its "full" signal goes. The Principal candidate notices that the signal crosses team boundaries — the producer is another org's service — so the backpressure protocol is a contract, not an implementation detail.

Why "Queue sizing" separates levels

A queue sized in items hides its real cost. 10,000 slots at 2,000 req/s is a 5-second wait; with a 2-second client timeout, every item past slot 4,000 is processed for nobody. The Staff answer sizes the queue to the latency budget and treats anything older than the deadline as already dropped. The Principal answer makes that the default in the shared library, because the 40th team will otherwise ship new LinkedBlockingQueue() with no bound.

Why "Failure" separates levels

The signature of missing backpressure is a system that is 100% busy and 0% useful: CPU pegged, throughput normal, success rate collapsing because every response arrives after the caller gave up. Seniors read the CPU graph and add instances. Staff engineers read queue.oldest_item_age_ms against the client deadline and see that the fix is to drop, not to add.


Where the Design Splits#

#Fault LineThe Tension
1Push vs PullProducer-driven speed (low latency, needs a stop signal) vs consumer-driven demand (safe by construction, adds a round trip)
2Buffer vs RejectAbsorb the burst (lossless, adds latency) vs fail fast (protects goodput, visible errors)
3Static vs Adaptive LimitsPredictable, hand-tuned thresholds vs limits that track real capacity and sometimes guess wrong
4Where to Push BackClient, gateway, queue, or service — each hop that absorbs the signal shields everything above it

Fault Line 1: Push vs Pull#

In push, the producer sends as fast as it can and the consumer must say "stop": TCP's zero receive window (Latency & Protocols), an HTTP 429, a WebSocket server closing a slow client. In pull, the consumer asks for work it has room for: a Kafka consumer's poll() returns at most max.poll.records (default 500), a Reactive Streams subscriber calls request(n), an SQS worker receives up to 10 messages. Credit-based flow control (HTTP/2 WINDOW_UPDATE, Flink's network credits) is the hybrid — the producer pushes, but only up to credits the consumer granted. Who pays: push without a stop signal makes the consumer pay with memory; pull makes the system pay one round trip of latency per batch and requires the consumer to be honest about capacity. Staff default: pull or credits for any hop crossing a process boundary in an async pipeline; push only with an explicit, enforced stop signal. Deviate when: latency is the product (market data, game state) — push, and drop on overflow rather than backpressure.

Fault Line 2: Buffer vs Reject#

A buffer is only useful while the work in it is still worth doing. For a request with a 1-second client timeout, a 5-second queue is a waste machine. For a billing event with a 24-hour SLA, a 6-hour backlog is fine. Who pays: buffering pays in latency and stale work; rejecting pays in visible errors and depends on the caller retrying sensibly. Staff default: buffer up to the latency budget measured in time, reject beyond it, and drop items whose deadline has passed at dequeue (check age before doing work). Deviate when: the work is lossless by contract — then buffer durably (a log, not memory) and make lag the alerting signal.

Fault Line 3: Static vs Adaptive Limits#

A static "max 200 concurrent" is correct on the day it was load-tested. After a deploy that adds 15ms, a cold cache, or a move to a smaller instance, it's wrong in one direction or the other. Adaptive limiters infer capacity from latency the way TCP congestion control infers bandwidth from loss and delay: when latency rises above the no-load baseline, shrink the limit; when it holds, probe upward. Who pays: static limits pay with either needless rejections or a collapsed backend; adaptive limits pay with oscillation and some false rejections while probing. Staff default: adaptive concurrency limit (gradient- or AIMD-style) on the server, with a static ceiling as a safety rail and a static floor so it can't collapse to zero. Deviate when: the limit encodes a contract, not capacity — per-tenant quotas are static on purpose.

Fault Line 4: Where to Push Back#

Overload signals can be absorbed at four places, and the cheapest one is the furthest upstream that can afford to wait.

WhereMechanismAbsorbsCost
ClientRetry budget, adaptive throttling, jittered backoff, honouring Retry-AfterRetry amplification; avoids sending doomed requestsNeeds client cooperation; you don't control third-party clients
Gateway / edgeRate limits, concurrency caps, priority admissionAbuse and bursts before they consume backend resourcesCoarse: can't see per-backend capacity without feedback
Queue / logDurable buffering, consumer pull, quotas that delay producersMinutes to days of mismatchLag; retention as a hard deadline; storage
ServiceBounded work queue, adaptive limiter, deadline check at dequeueLast line; knows true capacityWork already paid for (network, auth, parsing) before rejection

Staff default: every layer does something, but the service-level limiter is mandatory because it is the only one that knows real capacity; the gateway enforces contracts; clients enforce retry budgets (Stopping Cascading Failures covers the client side in depth). Deviate when: a durable log sits in the path — then the service behind it pulls, and the edge never needs to know the consumer is slow.


Common Interview Mistakes#

What Candidates SayWhat Interviewers HearWhat Staff Engineers Say
"Put a queue in front of it to absorb spikes""A queue is capacity""A queue moves overload forward in time. I'll size it to 100ms of drain and reject past that, because a 2s client timeout makes anything older worthless."
"Autoscale the consumers on CPU""Hasn't watched a 4-minute scale-up lose to a 30-second burst""I'll scale on lag-seconds, keep 25–30% headroom so a one-hour stall drains in under four hours, and the log absorbs the burst meanwhile."
"Return 503 when overloaded""Rejection without a story for retries""503 with Retry-After, and clients have a 10% retry budget — otherwise every rejection comes back three times and we never recover."
"Set a rate limit of 5,000 RPS on the service""Static number that's wrong after the next deploy""Rate limits are tenant contracts at the gateway. Capacity protection is an adaptive concurrency limit in the service, derived from latency."
"Use backpressure so we never drop anything""Doesn't know where backpressure ends""Backpressure has to stop somewhere. Here it stops at the topic. Above that, the API sheds, because a user can't wait 10 minutes."
"Monitor queue depth""Count, not time""Alert on the age of the oldest item and on time-to-drain. Five million messages is either 50 seconds or 14 hours."

Quick Reference#

Diagram: Quick Reference

Staff Sentence Templates#

"The bottleneck is [stage]. When it's full, the signal travels to [upstream hop] via [credits / blocking / 429], and it stops at [durable log / the edge], where it turns into [lag / a fast rejection]. [Team] owns that boundary."

"I'm sizing this queue in time, not items: [drain rate] × [wait budget] = [N] slots. Anything older than [deadline] at dequeue is dropped, because the caller has already given up."

"Consumer lag is [M] messages at [R] msg/s, so [M/R] seconds behind. With [H]% headroom, a one-hour stall drains in [1/H] hours. Product signed off on [T] as the maximum time-to-process."

"This hop is lossless, so it backpressures. That hop is latency-bound, so it sheds. The difference is whether waiting still has value, and I want that written down per hop."


Implementation Deep Dive#

1. Time-Bounded Work Queue with Deadline Drop — Any RPC Server#

The default for a synchronous service. Two additions most candidates skip: size the queue by time, and check the request's remaining deadline at dequeue, not just at enqueue.

DRAIN_RATE      = 2000          # req/s this instance sustains at target latency
WAIT_BUDGET_MS  = 100           # max time a request may wait before work starts
QUEUE_CAP       = DRAIN_RATE * WAIT_BUDGET_MS / 1000   # = 200 slots

queue = BoundedQueue(capacity=QUEUE_CAP)

function onRequest(req):
    if req.deadline - now() < MIN_USEFUL_MS:                 # already doomed
        metrics.incr("admission.rejected", tags=["reason:deadline"])
        return reply(503, retry_after="1")
    if not queue.offer(req):                                 # never block the I/O thread
        metrics.incr("admission.rejected", tags=["reason:queue_full"])
        return reply(503, retry_after="1")

worker loop:
    req = queue.take()
    metrics.histogram("queue.wait_ms", now() - req.enqueued_at)
    if now() >= req.deadline:                                # expired while waiting
        metrics.incr("queue.expired_dropped")
        continue                                             # do not do dead work
    handle(req, deadline=req.deadline)                       # propagate remaining budget downstream

Why it matters: the Google SRE book recommends keeping queue lengths small relative to the thread pool (on the order of 50% or less for steady traffic) and propagating deadlines so downstream servers don't work on requests whose callers have gone. The dequeue check is the cheapest goodput fix in existence: one comparison per request.

🎯 Staff Insight: Under sustained overload, a FIFO queue serves the oldest requests first — the ones most likely to have timed out. Switching to LIFO once a queue forms (Facebook's "adaptive LIFO") gives the freshest requests the best chance of meeting their deadline. It's unfair to the old ones, but they were already lost.

Pull-based demand is the cleanest form of backpressure: the producer can never send more than the consumer asked for. The Reactive Streams specification makes this a hard rule — the number of onNext signals must never exceed the total requested via request(n).

# Consumer side: request only what fits in the local buffer
class BoundedSubscriber:
    BUFFER = 256
    LOW_WATER = 64

    onSubscribe(sub):
        self.sub = sub
        sub.request(BUFFER)                  # initial credit

    onNext(item):
        buffer.push(item)                    # never exceeds BUFFER by spec rule 1.1

    worker loop:
        process(buffer.pop())
        self.outstanding -= 1
        if self.outstanding <= LOW_WATER:    # replenish in batches, not 1-by-1
            self.sub.request(BUFFER - self.outstanding)
            self.outstanding = BUFFER
Diagram: 2. Credit-Based Flow Control — Reactive Streams and Flink

Why credits beat TCP alone: stream processors multiplex many logical channels over one TCP connection. If one slow receiver stops reading, TCP's window closes for all of them — one hot subtask stalls its healthy neighbours. Per-channel credits let the slow channel stop while the rest continue. The rule generalises: flow control must be at the granularity of the thing that can be slow. HTTP/2 does the same thing with per-stream and per-connection windows (initial 65,535 octets each).

3. Consumer Lag as the Backpressure Signal — Kafka#

With a durable log, the producer is never slowed by the consumer at all; the log absorbs the difference and lag is the backpressure. The consumer's job is to pull at a rate it can sustain and to stop pulling when its own downstream is slow.

consumer = KafkaConsumer(
    max_poll_records = 200,             # default 500; smaller = smaller blast radius per batch
    max_poll_interval_ms = 300000,      # default 5 min; exceed it and you're kicked from the group
    enable_auto_commit = false)

inflight = Semaphore(400)               # bounded handoff to async workers

loop:
    if inflight.available() < 100:
        consumer.pause(assigned_partitions)      # stop fetching, keep group membership
    elif consumer.paused():
        consumer.resume(assigned_partitions)
    records = consumer.poll(timeout=100ms)       # returns nothing while paused, still heartbeats
    for r in records:
        inflight.acquire()
        submit(r, on_done=commit_and_release)

every 15s:
    lag_msgs  = sum(end_offset(p) - committed(p) for p in partitions)
    lag_sec   = lag_msgs / max(consume_rate_1m, 1)
    drain_sec = lag_msgs / max(consume_rate_1m - produce_rate_1m, 1)
    metrics.gauge("consumer.lag_seconds", lag_sec)
    metrics.gauge("consumer.time_to_drain_seconds", drain_sec)

The trap: a consumer that calls poll() and then blocks for minutes on a slow database exceeds max.poll.interval.ms, gets evicted, triggers a rebalance, and the partition's next owner hits the same slow database — a rebalance storm that cuts throughput to zero. pause()/resume() is how a Kafka consumer applies backpressure to itself without leaving the group.

🎯 Staff Move: "I'll alert on lag in seconds and on time-to-drain, not on message count. If time-to-drain exceeds the retention window minus a safety margin, that's a page — we're about to lose data, not just be late."

4. Adaptive Concurrency Limit — Gradient Style#

Static limits encode last quarter's capacity. A gradient limiter tracks the no-load latency (minRTT) and shrinks the limit when measured latency rises above it — queueing is the signal.

limit     = 50                   # starting point
MIN, MAX  = 10, 1000             # safety rails
BUFFER    = 0.5                  # tolerate 50% latency over baseline before shrinking

every sample window (e.g. 250 requests or 1s):
    sample_rtt = p50(latencies in window)
    if time_to_remeasure_min_rtt():          # periodically, with jitter across hosts
        pin limit to MIN briefly; min_rtt = observed p50
    gradient  = clamp((min_rtt * (1 + BUFFER)) / sample_rtt, 0.5, 1.0)
    headroom  = sqrt(limit)                  # lets it probe upward
    limit     = clamp(gradient * limit + headroom, MIN, MAX)

onRequest(req):
    if inflight >= limit:
        metrics.incr("limiter.rejected")
        return reply(503)                    # shed immediately, don't queue
    inflight += 1
    try: return handle(req)
    finally: inflight -= 1; record(latency)

This is the shape of the gradient controller in Envoy's adaptive concurrency filter (gradient = (minRTT + B) / sampleRTT, headroom = square root of the limit) and of Netflix's open-source concurrency-limits library. Envoy's docs warn that the periodic minRTT re-measurement temporarily drops the limit and can produce a visible burst of 503s — jitter it across hosts.

🎯 Staff Insight: An adaptive limiter can't tell its own overload from a slow dependency. If the database gets slow, latency rises, the limiter shrinks, and you shed traffic that would have succeeded slowly. That's usually right — it stops you piling more onto the database — but say it out loud, and make sure the dashboard shows which latency moved.

Technique Comparison

TechniqueSignalLatency of SignalLossless?Needs Producer CooperationBest For
Bounded queue + rejectQueue full / deadlineImmediateNo (sheds)NoSync RPC servers
Blocking putThread blocksImmediateYesImplicitlyIn-process pipelines, batch
Credit / request(n)Credits withheld1 round tripYesYes (protocol)Stream processing, reactive I/O
Durable log + pullLag growsMinutes (by design)Yes, until retentionNoAsync event pipelines
Broker quota (delay)Throttled responsePer requestYesClient honours delayMulti-tenant brokers
Adaptive concurrencyLatency over baselineSecondsNo (sheds)NoService self-protection
Retry budget / client throttlingRecent rejection ratioSecondsNoYes (client lib)Stopping amplification

Architecture Diagram#

Diagram: Architecture Diagram

How to narrate it: the dashed edges are backpressure. Warehouse slowness travels back to the fulfilment consumers, which pause, so the topic absorbs it as lag — and stops there. Nothing to the left of the topic ever knows the warehouse is slow. To the left of the topic, the signal is a fast 503 or 429, because a person is waiting. The topic is where backpressure turns from "wait" into "be late," and its retention is the hard deadline on that lateness.


Failure Scenarios#

1. The Busy-but-Useless Service — 100% CPU, 3% Goodput#

A search service normally runs at 60% CPU with 1,500 req/s and p99 of 80ms. Each instance has a 64-thread pool and an unbounded request queue. Clients time out at 1 second and retry up to 3 times.

t=0      Index refresh doubles per-query cost. Capacity drops from ~2,500 to ~1,250 req/s per cluster.
t=+5s    Queue grows by ~250 req/s cluster-wide. Queue wait passes 1s (1,250 items at 1,250/s).
t=+10s   Clients time out and retry. Offered load 1,500 -> ~4,000 req/s.
t=+60s   Queue wait 8s. Every dequeued request has already timed out. CPU 100%.
t=+90s   Success rate 3%. Health checks queue behind real traffic, fail. LB ejects hosts.
t=+3min  Fewer hosts, same offered load. Cluster-wide collapse.
t=+11min On-call restarts with traffic blocked at the gateway, then ramps 10% per minute.
t=+20min Refresh finished long ago. The overload outlived its cause by 18 minutes.

Detection: queue.oldest_item_age_ms > client timeout; queue.expired_dropped (if it existed — it didn't); ratio of requests.completed_after_deadline to completed; retry ratio > 10%. Blast radius: all search traffic; health-check failure turned partial overload into host ejection. Mitigation: cap queue at ~100ms of drain, drop expired at dequeue, adaptive concurrency limit, 503 instead of hang. Prevention: client retry budget (10%), health checks served on a separate, unqueued path, overload game day. Owner: search team owns the limiter; the RPC platform owns the queue defaults and the health-check path.

🎯 Staff Insight: This is a metastable failure: the trigger (index refresh) ended at t=+2min, but the retries and the queue of dead work kept the system overloaded. The fix that matters is not more capacity; it's making sure the system sheds fast enough that overload can't sustain itself.

2. Lag That Outran Retention — 1.7 Billion Events Lost#

A clickstream pipeline produces 80K events/s into a topic with 24h retention. Consumers process 100K events/s at full scale, so headroom is 25%.

Fri 18:00  Consumer deploy with a bug: one malformed event throws, consumer retries it in-line forever.
Fri 18:05  Consume rate 0. Lag alert (threshold 10M messages) fires, routes to a low-urgency queue.
Sat 18:00  Oldest unconsumed events reach 24h and begin aging out of retention.
Sun 00:00  Bug found, poison event skipped. 6h of events (~1.7B) already deleted.
Sun 00:00  Backlog = 24h x 80K/s = ~6.9B. Net drain = 100K - 80K = 20K/s.
Sun 00:00  Time to drain = 6.9B / 20K = ~96h. Dashboards are 4 days late on top of the loss.

Detection: consumer.lag_seconds > 900 (page); consumer.time_to_drain_seconds > retention − elapsed lag − 6h (page: data loss imminent); poison-message counter; consume rate = 0 for 5 min (page). Blast radius: 6 hours of analytics data permanently lost; all downstream reporting late for ~4 days. Mitigation: dead-letter the poison event instead of retrying in-line; extend topic retention the moment lag-seconds crosses half of it; burst-scale consumers 3× (partition count allowing) to raise net drain from 20K/s to ~220K/s, cutting a 96-hour drain to under 9. Prevention: alert on lag in seconds and time-to-drain, not messages; partition count sized for 3× burst consumption; DLQ with an owner. Owner: pipeline team owns consumers and the DLQ; data platform owns retention and the time-to-drain alert template.

3. One Slow Tenant Stalls the Shared Queue#

A notification platform uses a single shared work queue for all tenants. One tenant sends 2M pushes in 3 minutes; their provider starts returning 429s and the workers back off in place, holding queue slots while they sleep.

t=0      Tenant A enqueues 2M messages; normal total is 3K/s.
t=+30s   Provider throttles tenant A. Workers retry with backoff, holding the work item.
t=+2min  90% of workers are sleeping on tenant A. Tenant B's password-reset emails wait 6 min.
t=+8min  On-call manually purges tenant A's messages. Tenant A's campaign is lost.

Detection: queue.wait_ms p99 by tenant; worker.sleeping_ratio; per-tenant share of in-flight work > 50%. Blast radius: every tenant on the shared queue, including critical transactional mail. Mitigation: per-tenant queues (or shuffle-sharded queues), per-tenant concurrency caps, and backoff that releases the slot and re-enqueues with a delay. Prevention: never let a worker sleep while holding shared capacity; backpressure from one downstream must only slow the tenants that use it. Owner: notification platform owns isolation; tenant A's team owns their campaign pacing.

Operational Reality Matrix#

FailureDetection SignalBlast RadiusMitigationOwner
Unbounded queue of dead workqueue.oldest_item_age_ms > client timeoutWhole service; goodput near 0Time-bounded queue, drop expired, LIFO under loadService team + RPC platform
Retry amplificationclient.retry_ratio > 10%Every backend in the chainRetry budgets, adaptive client throttlingClient library owners
Consumer lag beyond retentionconsumer.time_to_drain_seconds vs retentionPermanent data lossDLQ poison items, burst-scale consumers, extend retentionPipeline team
Rebalance storm from slow handlerconsumer.rebalances_per_min > 2Consumer group throughput to 0pause()/resume(), smaller max.poll.recordsPipeline team
Head-of-line blocking across tenantsPer-tenant queue.wait_ms divergenceAll tenants on shared queuePer-tenant or shuffle-sharded queues, release-on-backoffPlatform team
Limiter oscillationlimiter.limit swings > 50% per minuteFalse 503sSmooth with longer windows, raise floor, jitter minRTT probesService team
Blocking propagates into unrelated workThread pool saturation on a shared executorUnrelated endpointsBulkhead executors per downstreamService team

Beyond Staff: The Principal View#

Why L7 Sees This Problem Differently#

A Staff engineer makes one pipeline or one service degrade gracefully. A Principal engineer notices that the org's worst incidents of the last year were not caused by any single overloaded service but by the seams: a team whose client retried 5× against a team whose server queued without bound, a shared broker where one tenant's burst became everyone's lag, a gateway that returned 503 to clients that treated 503 as "retry immediately." Backpressure only works if every hop speaks the same protocol, and at org scale that protocol is a set of cross-team contracts: what a rejection looks like, what a caller must do when it gets one, how much headroom a consumer must keep, and where the org has decided waiting ends and dropping begins. At L7 the work is making those contracts the default in shared libraries so no team has to rediscover them during an incident.

🧭 Principal Move: "Our last three Sev-1s were overload that one team's retries turned into another team's outage. I'd standardize the overload contract — one rejection format, mandatory retry budgets and Retry-After in the shared clients, time-bounded queues in the RPC framework — so backpressure works across team boundaries by default, not by heroics."

The Org-Level Fault Line#

One overload protocol in shared infrastructure vs per-team overload handling.

OptionWhat WorksWhat BreaksWho Pays
Per-team limiters and retry logicAutonomy; each team tunes for its own trafficIncompatible semantics at every seam: one team's 503 is another's "retry now"; amplification across 4–5 hopsOn-call of whichever team is at the bottom of the call graph
Central "overload service" that every call consultsOne view of capacityAdds a network hop and a single point of failure to the hot path, exactly when it's overloadedEveryone, simultaneously
Protocol in shared libraries + mesh, policy per teamRetry budgets, deadline propagation, time-bounded queues and adaptive limits on by default; teams only set criticality and quotasPlatform team must own and evolve RPC, queue and Kafka client libraries; exotic needs waitPlatform headcount (3–6 engineers)

The Principal default is the third row: mechanism is centralised in code everyone links or every sidecar runs; policy — which requests are sheddable, which pipelines are lossless, what each tenant's quota is — stays with the owning teams and is declared, not hard-coded.

Cost Model#

Assumptions: fully loaded engineer ~$25K/month; compute at ~$0.05 per vCPU-hour; Kafka storage ~$0.10/GB-month per copy with 3× replication (so ~$0.30 per logical GB-month); "headroom" means extra consumer capacity beyond peak produce rate to drain backlog.

ScaleTrafficOverload MachineryHeadroom CostPeople / On-callRough Monthly Total
Startup2K req/s, 5K events/sStatic limits, bounded queues, one Kafka cluster20% consumer headroom ≈ 4 vCPU (~$150)0.2 FTE; service teams on-call~$1K + ~$5K people
Growth50K req/s, 300K events/sAdaptive limiters in the RPC library, retry budgets, lag-seconds alerting, 3-day retention30% headroom ≈ 120 vCPU (~$4.5K); extra 2 days retention on 25 TB/day ≈ $15K1.5 FTE platform; shared rotation~$25K + ~$40K people
Large1M req/s, 5M events/sMesh-enforced limits and deadlines, per-tenant quotas, shuffle-sharded queues, overload game days30% headroom ≈ 2,000 vCPU (~$75K); 7-day retention on 400 TB/day ≈ $800K5-person platform team; dedicated rotation~$900K + ~$125K people

The Principal observation: at the large scale, retention and drain headroom dominate — and both are insurance, not throughput. Seven days of retention exists so a weekend-long consumer bug doesn't become data loss; 30% consumer headroom exists so a one-hour stall drains in ~3.3 hours instead of 20. Both look like waste in a cost review. They should be explicit budget lines with a stated recovery objective ("any 24h consumer outage drains within 72h, with zero loss"), owned by the data platform, so nobody trims them without trimming the promise.

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

The Year 3 trigger is the one teams miss: once dozens of services pin old client library versions, retry budgets and deadline propagation silently disappear from part of the call graph. Moving enforcement into the mesh (or a mandatory sidecar) is how the contract survives version skew.

One-Way Doors vs Two-Way Doors#

DecisionDoor TypeReversibility Cost
Queue capacity, limiter bounds, retry budget percentageTwo-wayConfig change, minutes
Consumer prefetch / max.poll.recordsTwo-wayDeploy, minutes
Kafka partition count for a topic (caps consumer parallelism, so caps drain rate)One-way-ishIncreasing is possible but re-maps keys and breaks per-key ordering during the change; plan for 3× burst consumption upfront
Making an API synchronous that could have been async (202 + log)One-wayClients build around the response; moving to async later is a breaking contract change
Publishing a public quota or rejection contract (429 semantics, Retry-After)One-wayThird-party clients encode it; changing semantics breaks integrations
Declaring a pipeline losslessOne-way-ishDownstream teams build reconciliation-free systems on it; relaxing it later needs every consumer's sign-off
Choosing a pull-based vs push-based transport between two teamsTwo-way (with cost)One quarter of dual-running and client migration

The Standard I'd Write#

RFC: Overload and Backpressure Contract (v1)

Scope: Every service-to-service RPC, every queue or topic consumed by more than one team, and every public API.

MUST:

  1. In-memory queues in serving paths MUST be bounded, with capacity derived from a declared wait budget (default 100ms). Items past their deadline MUST be dropped at dequeue and counted.
  2. Rejections due to overload MUST use 503 (or gRPC UNAVAILABLE / RESOURCE_EXHAUSTED) with a Retry-After hint. Callers MUST honour it.
  3. Shared clients MUST enforce a per-client retry budget (default 10% of requests) and propagate deadlines.
  4. Every async pipeline MUST declare itself lossless or sheddable, publish consumer.lag_seconds, and page on time-to-drain exceeding 50% of retention.
  5. Consumers of shared topics MUST keep enough partitions and headroom to drain a 1-hour stall within 4 hours.

SHOULD: Use the platform adaptive concurrency limiter; run one overload game day per quarter per tier-1 service; isolate tenants with per-tenant or shuffle-sharded queues.

Exceptions: Filed with the platform team, time-boxed to 2 quarters, with an owner and a stated alternative protection.

Success metrics: zero Sev-1s where retries sustained an overload after the trigger ended; p99 queue.oldest_item_age_ms below the declared budget on 95% of services; zero retention-expiry data loss on lossless pipelines.

What I'd Tell the VP#

When traffic spikes or something downstream slows, our systems currently pile the extra work into queues and retries until everything falls over — and stays down after the spike is gone. That's what turned two 2-minute blips last year into 20-minute outages. I'm proposing a single set of rules, built into the shared libraries every team already uses, so that busy systems say "not now" quickly and callers back off politely. It costs about two engineers for two quarters and some extra storage and spare capacity that we'll budget explicitly as insurance. In exchange, short spikes stay short, and we stop losing event data when a pipeline breaks over a weekend.

Principal Interview Signals#

SignalWhat It Sounds Like
Contracts at the seams"The overload bug lives between teams. I'd standardise what a rejection means and what callers must do with it."
Pricing the insurance"30% consumer headroom and 7-day retention are our recovery objective in dollars. They go on a budget line with the promise attached."
Mechanism vs policy"The platform ships the limiter and the retry budget; teams declare what's sheddable. I don't want 40 hand-tuned retry loops."
One-way door awareness"Partition count caps drain rate forever-ish, so I size it for 3× burst consumption now."
Rehearsal over hope"We don't know where backpressure ends until we game-day it. Quarterly, per tier-1 service, with the drop boundary written down first."

Staff answers that L7 interviewers find insufficient:

  • "I'd add an adaptive limiter and a retry budget to this service." — Right for one service; silent on the 30 clients that still retry 5× and the library they'll copy from.
  • "We'll keep enough consumer headroom to catch up." — Doesn't say how much, what it costs, or who defends it in the next cost review.
  • "Backpressure all the way to the client." — Doesn't decide, as an org, where waiting ends and dropping begins — so every team decides differently.

How Real Companies Built It#

Flink's network stack has used credit-based flow control since Flink 1.5. Each input channel gets a small number of exclusive network buffers (default 2 per channel) plus a shared pool of floating buffers per input gate (default 8), with 32 KB memory segments by default. Receivers announce free buffers to senders as credits (one buffer equals one credit), and senders announce their backlog so receivers can request floating buffers. The motivation was multiplexing: many logical channels share one TCP connection, so before credits a single slow receiver could stall every other subtask on that connection. Flink 1.14 added automatic buffer debloating to shrink in-flight data, since buffered data also lengthens checkpoint alignment (Flink network stack deep dive, Flink network memory tuning).

Staff insight: This is the cleanest real example of "flow control must be at the granularity of the thing that can be slow." TCP gave Flink backpressure for free, at the wrong granularity. It also shows the hidden cost of buffers: every buffer you add for throughput is data that a checkpoint has to wait for or persist. See Flink.

Netflix and Envoy: Adaptive Concurrency Limits#

Netflix's open-source concurrency-limits library replaces hand-set RPS limits with a concurrency limit derived from Little's law and adjusted from observed latency, borrowing from TCP congestion control (Vegas- and gradient-style algorithms, plus AIMD for clients). When the limit is reached, excess requests are rejected immediately — HTTP 429 from the servlet filter or UNAVAILABLE in gRPC — instead of being queued. Envoy ships the same idea as a proxy filter: a gradient controller computes (minRTT + buffer) / sampleRTT, adds square-root headroom, and periodically re-measures minRTT by pinning the limit low, which its docs warn can cause a visible increase in 503s during the measurement window (concurrency-limits, Envoy adaptive concurrency).

Staff insight: Two independent implementations converged on "measure queueing via latency, shed the excess." Citing both lets you defend adaptive limits as a mainstream production technique and name its sharp edge — the probe that briefly sheds real traffic.

Facebook: CoDel and Adaptive LIFO for Server Queues#

Facebook's "Fail at Scale" article describes applying the controlled-delay idea from network queue management to server request queues: if a queue has not been empty for the last N milliseconds, time spent in the queue is capped at M milliseconds, with 5ms and 100ms reported as working well in practice. It pairs this with adaptive LIFO — FIFO normally, LIFO once a queue forms — so the newest requests, the ones whose callers are still waiting, are served first. The implementation lives in their open-source Wangle library used by Thrift (ACM Queue: Fail at Scale).

Staff insight: This is "buffer by time, not by count" turned into a mechanism. A queue that is allowed to stand for long periods is serving dead requests. Interviewers rarely expect the algorithm name; they do expect you to know FIFO is the wrong order under sustained overload.

Apache Kafka: Lag, Bounded Polls and Throttling Quotas#

Kafka shows two backpressure styles in one system. On the consumer side it is pull-based: poll() returns at most max.poll.records (default 500) and a consumer that doesn't poll within max.poll.interval.ms (default 5 minutes) is removed from the group, so slow handlers must pause fetching rather than block. On the broker side, client quotas throttle rather than error: a broker that detects a quota violation computes the delay needed to bring the client back under quota, returns a response carrying that delay, and mutes the client's channel until it passes; well-behaved clients also stop sending for that interval (Kafka consumer configs, Kafka quotas).

Staff insight: Quota-by-delay is backpressure aimed at a producer you share infrastructure with — it slows the noisy tenant without failing its writes. Pair it with lag-in-seconds alerting and you have both halves of a multi-tenant pipeline story. See Kafka and Message Broker.

Google SRE and the Reactive Streams Specification: The Rules Written Down#

Two primary sources state the rules this page relies on. The Google SRE book recommends small queues relative to thread pool size (about 50% or less for steady traffic), LIFO or CoDel to avoid serving requests that are already dead, deadline propagation, a per-request retry cap of 3 and a per-client retry budget of 10%, and warns that 3 layers each retrying 4 times can turn one user action into 64 database attempts. The Reactive Streams specification (adopted 1:1 as java.util.concurrent.Flow in JDK 9) makes demand-driven flow control a protocol rule: a publisher must never signal more onNext than the subscriber has requested (SRE: Addressing Cascading Failures, SRE: Handling Overload, Reactive Streams spec).

Staff insight: Quoting "64×" is the fastest way to make an interviewer take retry budgets seriously, and quoting Reactive Streams rule 1.1 is the fastest way to explain why pull-based systems can't be flooded.


Practice Drill#

Prompt: "Our order pipeline takes checkout events from the API, writes them to a queue, and a fulfilment service calls a warehouse partner API capped at 500 requests/sec. On sale days we take 3,000 orders/sec for 20 minutes. Last sale, the API timed out for 15 minutes and the warehouse partner blocked us for exceeding their limit. Redesign the overload handling."

Staff Answer

The two symptoms are the same bug: no hop knew where backpressure should stop. I'd draw the line at a durable log. The checkout API writes the order to Kafka with acks=all and returns 202 — that write takes ~5–10ms and is the only thing a buyer waits on. If the Kafka write itself exceeds 200ms, the API sheds with 503 and Retry-After rather than holding connections; that's the only overload a user ever sees. Everything to the right of the topic is lossless and late, not failed. The fulfilment consumers pull at most 200 records per poll, push into a bounded in-flight pool of ~400 calls, and call the warehouse through a client-side token bucket at 450/s — 10% under the contract so we never trip their ban. When the pool fills or the partner returns 429, consumers pause() their partitions instead of sleeping inside a handler, so they keep group membership and the topic absorbs the difference as lag. Sale-day math: 3,000/s for 20 minutes is 3.6M orders; at 450/s net drain (normal off-sale input is small) that's ~2.2 hours to fulfil. Product has to sign off on "orders confirmed instantly, sent to the warehouse within 3 hours on sale days" — or we negotiate a higher partner limit, which is a business conversation, not an engineering one. Topic retention 7 days, 24 partitions (enough for future 3× partner capacity). Metrics: consumer.lag_seconds, consumer.time_to_drain_seconds, warehouse.client.throttled, api.kafka_write_p99_ms, admission.rejected.

Why this is L6:

  • Places the backpressure boundary explicitly at the log and turns the user-facing edge into fast shedding, not waiting.
  • Uses pull + pause/resume + a client-side rate limiter so a hard downstream contract becomes lag, never errors or bans.
  • Converts the burst into a time-to-drain number and makes it a product decision with an owner.

What L7 adds:

  • Notices that every partner integration (payments, shipping, tax) has the same shape and proposes a shared "partner egress" service with per-partner token buckets and lag SLOs, instead of each team rebuilding this.
  • Prices it: the 2.2-hour backlog vs a partner contract upgrade, and the cost of 7-day retention as the insurance for a weekend-long consumer failure.
  • Writes the contract: partner limits live in one config owned by the integrations team; fulfilment owns lag SLOs; product owns the "within 3 hours" promise.

Staff Interview Application#

How to Introduce This Pattern#

"Before I add queues, I want to find the slowest stage and decide what happens when it's full. For each hop, either the producer waits, the producer slows, we buffer for a bounded time, or we drop. Here, I'll let backpressure flow up to the Kafka topic, where it becomes lag, and above that the API sheds fast, because a user can't wait more than a second."

Lead with the bottleneck, then the signal path, then the boundary where waiting turns into dropping, then who owns that boundary.

When NOT to Use This Pattern#

  • Producer and consumer already match with large margin: a service at 20% utilisation with a single caller doesn't need adaptive limiters or credit protocols — a bounded queue and a timeout are enough. Add machinery when the measured peak-to-capacity ratio exceeds ~0.7.
  • Data where only the latest value matters: sensor gauges, presence, cursor positions, live video frames. Backpressure makes everything late; drop-oldest or sample instead.
  • Producers you cannot slow: third-party webhooks and device fleets ignore your stop signals. Accept into a durable log and treat it as a lag problem, or enforce a published quota at the edge.
  • Overload is actually abuse: a single tenant or attacker generating the load needs Rate Limiter per identity, not system-wide backpressure that slows every honest caller too.
  • Shedding by business priority: deciding which users get degraded is a separate design — see Degraded Mode Framework.

Follow-Up Questions to Anticipate#

Interviewer AsksWhat They Are TestingHow to Respond
"What's the difference between backpressure and load shedding?"Conceptual precision"Backpressure slows the producer; shedding drops the work. Backpressure is right while someone upstream can afford to wait; at the first hop where nobody can, it has to become shedding."
"How big should the queue be?"Time-based sizing"Drain rate times the wait budget. At 2,000/s and 100ms, 200 slots. Anything older than the caller's deadline is dropped at dequeue."
"The consumer is falling behind. What do you do?"Lag in time, headroom math"Convert lag to seconds and time-to-drain. If time-to-drain approaches retention, extend retention and burst-scale up to the partition count; then fix the root cause."
"Why not just autoscale?"Timescales"Autoscaling takes minutes; bursts take seconds; and the bottleneck is often a database or partner that doesn't scale. Autoscaling sets steady-state capacity; backpressure handles the gap."
"How do retries interact with this?"Amplification"Every rejection can come back multiplied. Per-client retry budgets of ~10% and honouring Retry-After keep a rejection from becoming three."
"Push or pull for this pipeline?"Flow-control models"Pull, or credits, for anything lossless between processes — the consumer can't be flooded by construction. Push with drop-on-overflow for latency-critical, superseding data."

Scorecard#

DimensionSenior (L5)Staff (L6)Principal (L7)
FramingAdds a queue and autoscalingFinds the bottleneck and traces the overload signal hop by hopTreats overload semantics as a cross-team contract
Backpressure vs sheddingUses the words interchangeablyPlaces the boundary explicitly and justifies it per hopSets org policy for which tiers wait and which drop
SizingQueue length in itemsQueue in time; lag in seconds; drain-time headroom mathRetention and headroom as funded recovery objectives
LimitsStatic RPS per endpointAdaptive concurrency + retry budgets + deadlinesMechanism in shared libraries/mesh, policy declared per team
OwnershipImplicitNamed owners for lag SLOs, DLQs and drop decisionsExceptions process and game days across tier-1 services

Strong Hire Signals

SignalWhat It Sounds Like
Distinguishes the tools"This hop backpressures, that hop sheds, and here's why."
Time-based sizing"200 slots is 100ms at this drain rate; beyond that the caller is gone."
Lag as time"Five million messages at 100/s is 14 hours, so this is a page."
Retry awareness"Without a budget, every 503 comes back three times."

Lean No-Hire Signals

SignalWhy It Misses the Bar
"Unbounded queue so we never lose anything"Will ship a busy-but-useless failure mode and an OOM
No answer for where the overload endsBackpressure with no terminal boundary just moves the outage
Lag alerts in message counts onlyCan't tell 50 seconds from 14 hours of delay

Common False Positives: Knowing Reactive Streams operator names ≠ knowing where backpressure should stop. Naming "token bucket" ≠ protecting a service's real capacity. Proposing Kafka ≠ having a plan for lag beyond retention.


Capacity Planning Quick Reference#

Sizing Formulas#

in_flight          = throughput × latency                      # Little's law
queue_capacity     = drain_rate × wait_budget                  # size in time
lag_seconds        = lag_messages / consume_rate
drain_time         = backlog / (consume_rate − produce_rate)
outage_drain_time  = outage_duration × produce_rate / (consume_rate − produce_rate)
                   = outage_duration / headroom_fraction       # 1h outage, 20% headroom → 5h
max_drain_rate     = partitions × per_consumer_rate            # partition count caps catch-up
retry_amplification (no budget) = attempts ^ layers             # 4 attempts × 3 layers → 64×

Key Numbers Worth Memorizing#

NumberContext
≤ 100–200 msQueue wait budget for a synchronous request path
~50% of pool sizeSRE guidance for queue length relative to thread pool, steady traffic
10%Per-client retry budget; 3 attempts per request cap
64×Load from 4 attempts at each of 3 layers
500 / 5 minKafka max.poll.records / max.poll.interval.ms defaults
65,535 octetsHTTP/2 initial per-stream and per-connection flow-control window
2 + 8Flink default exclusive buffers per channel + floating buffers per gate
5 ms / 100 msCoDel-style queue timeout / standing-queue interval reported by Facebook
1h ÷ headroomDrain time for a 1-hour stall: 5h at 20%, 3.3h at 30%, 20h at 5%
1–5 minTypical autoscaling reaction time — longer than most bursts

Common Pitfalls Checklist#

  • Every in-memory queue has a bound derived from a time budget
  • Requests past their deadline are dropped at dequeue, not processed
  • Each hop has a written answer: block, slow, buffer, or drop
  • Backpressure terminates at a durable log or a fast rejection — never an unbounded buffer
  • Clients have retry budgets and honour Retry-After
  • Lag is alerted in seconds and time-to-drain, with retention as the hard limit
  • Partition count allows at least 3× burst consumption for catch-up
  • Slow handlers pause fetching instead of blocking inside the poll loop
  • Shared queues isolate tenants; backoff releases the slot instead of sleeping on it
  • The overload boundary has been game-dayed, not just drawn
  1. Loading the index…