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
nitems 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.
| Intent | Constraint | Strategy | Failure Mode | Correctness 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 full | Lag grows faster than you can drain; retention expires before consumption | Zero loss; bounded time-to-process, product-signed |
| Latency-bound request path (APIs, RPC fan-out) | Caller waits ≤ 1–2 s; stale work is worthless | Bounded queues by time, adaptive concurrency limits, fast rejection, deadline propagation | Queues fill with expired requests; goodput collapses while CPU is pegged | Bounded p99; explicit 429/503 rather than timeouts |
| Fair shared infrastructure (multi-tenant brokers, platform APIs) | One tenant's burst must not slow everyone | Per-tenant quotas that throttle (delay) rather than error; separate queues per tenant class | Noisy neighbour fills the shared queue; head-of-line blocking | Per-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.
| Response | Mechanism | Cost Lands On | Right When |
|---|---|---|---|
| Block (synchronous backpressure) | Bounded queue put() waits; TCP zero window; Kafka producer max.block.ms | The producer's thread and its callers | Producer is a batch job or a durable log that can wait |
| Slow / throttle (signalled backpressure) | Credits, request(n), quota delays, 429 + Retry-After | Producer throughput, gracefully | Producer is cooperative and can pace itself |
| Buffer | Queue, log, spillover storage | Latency and memory; work goes stale | Burst is short relative to buffer-in-time |
| Drop (shedding) | Reject at admission, drop oldest, drop lowest priority | A user or a data consumer, visibly | Waiting has no value: the deadline has passed or the item is superseded |
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#
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Unbounded buffer | Never rejects; simple | OOM or hours of stale work; latency unbounded; overload outlives the spike | On-call at 3am; users waiting on work already abandoned |
| Bounded queue + block | Lossless; producer naturally paced | Blocking propagates to threads that serve unrelated work; deadlocks in cycles | Unrelated callers sharing the blocked thread pool |
| Credit / pull-based flow control | Producer cannot outrun consumer by construction; per-channel isolation | Needs protocol support end to end; one non-compliant hop breaks it | Platform team that owns the protocol and client libraries |
| Static concurrency or rate limit | Predictable, easy to reason about | Wrong the day capacity changes (deploy, cache cold, instance type) | Users rejected when there was capacity; or the backend when there wasn't |
| Adaptive concurrency limit | Tracks real capacity from latency; no hand-tuning | Oscillation, slow convergence, mistakes latency of a dependency for own overload | Some false rejections during probing; on-call debugging "why 503?" |
| Durable log as shock absorber | Decouples producer from consumer entirely; hours of slack | Lag becomes the failure mode; retention is a hard deadline; data goes stale silently | Data consumers seeing late results; storage bill |
| Shed at admission | Protects goodput; fast honest failure | Users see errors; needs priority to avoid dropping the wrong thing | Users, 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#
| Behavior | Senior (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 sizing | Large queue so nothing is lost | Size by time: drain rate × wait budget; reject beyond it; alert on oldest-item age | Sets an org-wide rule: no unbounded in-memory queues in serving paths; enforced in the shared RPC/queue libraries |
| Limits | Static RPS limit per endpoint | Adaptive concurrency limit from latency + retry budgets so rejection isn't undone | Funds the limiter as a platform default; prices false-rejection rate against outage minutes avoided |
| Async lag | Alert when lag > 1M messages | Alert on lag-seconds and time-to-drain; plan consumer headroom so a 1h outage drains in < 4h | Treats 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 |
| Ownership | Whoever owns the consumer | Producer owns pacing contract; consumer owns lag SLO; on-call owns drop decisions | Writes 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 Line | The Tension |
|---|---|---|
| 1 | Push vs Pull | Producer-driven speed (low latency, needs a stop signal) vs consumer-driven demand (safe by construction, adds a round trip) |
| 2 | Buffer vs Reject | Absorb the burst (lossless, adds latency) vs fail fast (protects goodput, visible errors) |
| 3 | Static vs Adaptive Limits | Predictable, hand-tuned thresholds vs limits that track real capacity and sometimes guess wrong |
| 4 | Where to Push Back | Client, 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.
| Where | Mechanism | Absorbs | Cost |
|---|---|---|---|
| Client | Retry budget, adaptive throttling, jittered backoff, honouring Retry-After | Retry amplification; avoids sending doomed requests | Needs client cooperation; you don't control third-party clients |
| Gateway / edge | Rate limits, concurrency caps, priority admission | Abuse and bursts before they consume backend resources | Coarse: can't see per-backend capacity without feedback |
| Queue / log | Durable buffering, consumer pull, quotas that delay producers | Minutes to days of mismatch | Lag; retention as a hard deadline; storage |
| Service | Bounded work queue, adaptive limiter, deadline check at dequeue | Last line; knows true capacity | Work 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 Say | What Interviewers Hear | What 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#
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.
2. Credit-Based Flow Control — Reactive Streams and Flink#
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
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
| Technique | Signal | Latency of Signal | Lossless? | Needs Producer Cooperation | Best For |
|---|---|---|---|---|---|
| Bounded queue + reject | Queue full / deadline | Immediate | No (sheds) | No | Sync RPC servers |
| Blocking put | Thread blocks | Immediate | Yes | Implicitly | In-process pipelines, batch |
Credit / request(n) | Credits withheld | 1 round trip | Yes | Yes (protocol) | Stream processing, reactive I/O |
| Durable log + pull | Lag grows | Minutes (by design) | Yes, until retention | No | Async event pipelines |
| Broker quota (delay) | Throttled response | Per request | Yes | Client honours delay | Multi-tenant brokers |
| Adaptive concurrency | Latency over baseline | Seconds | No (sheds) | No | Service self-protection |
| Retry budget / client throttling | Recent rejection ratio | Seconds | No | Yes (client lib) | Stopping amplification |
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#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Unbounded queue of dead work | queue.oldest_item_age_ms > client timeout | Whole service; goodput near 0 | Time-bounded queue, drop expired, LIFO under load | Service team + RPC platform |
| Retry amplification | client.retry_ratio > 10% | Every backend in the chain | Retry budgets, adaptive client throttling | Client library owners |
| Consumer lag beyond retention | consumer.time_to_drain_seconds vs retention | Permanent data loss | DLQ poison items, burst-scale consumers, extend retention | Pipeline team |
| Rebalance storm from slow handler | consumer.rebalances_per_min > 2 | Consumer group throughput to 0 | pause()/resume(), smaller max.poll.records | Pipeline team |
| Head-of-line blocking across tenants | Per-tenant queue.wait_ms divergence | All tenants on shared queue | Per-tenant or shuffle-sharded queues, release-on-backoff | Platform team |
| Limiter oscillation | limiter.limit swings > 50% per minute | False 503s | Smooth with longer windows, raise floor, jitter minRTT probes | Service team |
| Blocking propagates into unrelated work | Thread pool saturation on a shared executor | Unrelated endpoints | Bulkhead executors per downstream | Service 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.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Per-team limiters and retry logic | Autonomy; each team tunes for its own traffic | Incompatible semantics at every seam: one team's 503 is another's "retry now"; amplification across 4–5 hops | On-call of whichever team is at the bottom of the call graph |
| Central "overload service" that every call consults | One view of capacity | Adds a network hop and a single point of failure to the hot path, exactly when it's overloaded | Everyone, simultaneously |
| Protocol in shared libraries + mesh, policy per team | Retry budgets, deadline propagation, time-bounded queues and adaptive limits on by default; teams only set criticality and quotas | Platform team must own and evolve RPC, queue and Kafka client libraries; exotic needs wait | Platform 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.
| Scale | Traffic | Overload Machinery | Headroom Cost | People / On-call | Rough Monthly Total |
|---|---|---|---|---|---|
| Startup | 2K req/s, 5K events/s | Static limits, bounded queues, one Kafka cluster | 20% consumer headroom ≈ 4 vCPU (~$150) | 0.2 FTE; service teams on-call | ~$1K + ~$5K people |
| Growth | 50K req/s, 300K events/s | Adaptive limiters in the RPC library, retry budgets, lag-seconds alerting, 3-day retention | 30% headroom ≈ 120 vCPU (~$4.5K); extra 2 days retention on 25 TB/day ≈ $15K | 1.5 FTE platform; shared rotation | ~$25K + ~$40K people |
| Large | 1M req/s, 5M events/s | Mesh-enforced limits and deadlines, per-tenant quotas, shuffle-sharded queues, overload game days | 30% headroom ≈ 2,000 vCPU (~$75K); 7-day retention on 400 TB/day ≈ $800K | 5-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#
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#
| Decision | Door Type | Reversibility Cost |
|---|---|---|
| Queue capacity, limiter bounds, retry budget percentage | Two-way | Config change, minutes |
Consumer prefetch / max.poll.records | Two-way | Deploy, minutes |
| Kafka partition count for a topic (caps consumer parallelism, so caps drain rate) | One-way-ish | Increasing 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-way | Clients 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-way | Third-party clients encode it; changing semantics breaks integrations |
| Declaring a pipeline lossless | One-way-ish | Downstream 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 teams | Two-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:
- 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.
- Rejections due to overload MUST use 503 (or gRPC
UNAVAILABLE/RESOURCE_EXHAUSTED) with a Retry-After hint. Callers MUST honour it.- Shared clients MUST enforce a per-client retry budget (default 10% of requests) and propagate deadlines.
- Every async pipeline MUST declare itself lossless or sheddable, publish
consumer.lag_seconds, and page on time-to-drain exceeding 50% of retention.- 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_msbelow 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#
| Signal | What 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#
Apache Flink: Credit-Based Flow Control Between Operators#
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 Asks | What They Are Testing | How 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#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Framing | Adds a queue and autoscaling | Finds the bottleneck and traces the overload signal hop by hop | Treats overload semantics as a cross-team contract |
| Backpressure vs shedding | Uses the words interchangeably | Places the boundary explicitly and justifies it per hop | Sets org policy for which tiers wait and which drop |
| Sizing | Queue length in items | Queue in time; lag in seconds; drain-time headroom math | Retention and headroom as funded recovery objectives |
| Limits | Static RPS per endpoint | Adaptive concurrency + retry budgets + deadlines | Mechanism in shared libraries/mesh, policy declared per team |
| Ownership | Implicit | Named owners for lag SLOs, DLQs and drop decisions | Exceptions process and game days across tier-1 services |
Strong Hire Signals
| Signal | What 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
| Signal | Why 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 ends | Backpressure with no terminal boundary just moves the outage |
| Lag alerts in message counts only | Can'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#
| Number | Context |
|---|---|
| ≤ 100–200 ms | Queue wait budget for a synchronous request path |
| ~50% of pool size | SRE 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 min | Kafka max.poll.records / max.poll.interval.ms defaults |
| 65,535 octets | HTTP/2 initial per-stream and per-connection flow-control window |
| 2 + 8 | Flink default exclusive buffers per channel + floating buffers per gate |
| 5 ms / 100 ms | CoDel-style queue timeout / standing-queue interval reported by Facebook |
| 1h ÷ headroom | Drain time for a 1-hour stall: 5h at 20%, 3.3h at 30%, 20h at 5% |
| 1–5 min | Typical 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