Hiring BarSupport

Design a Distributed Message Queue — Staff-Level Case Study

Case study89 min read7 diagrams

Technologies referenced in this case study: Apache Kafka · ZooKeeper & etcd · Redis · PostgreSQL · DynamoDB

Related: Idempotency · Stream Processing · Task Scheduling · Notification Systems · Scaling Writes · Degraded Mode · Long-Running Processes · Data Pipelines · Consistency Models · Sharding

How to Use This Case Study#

This case study is organized for the interview first and for reference second. Read it front-to-back once; afterwards, return to the fault lines and deep dives that expose your weak spots. Broker internals (segments, ISR, the controller) live in the Kafka guide — this page is about the design problem: what the queue promises, to whom, and who owns the messages nobody can process.

ModeTimeWhat to Read
Quick Review15 minExecutive Summary → Interview Walkthrough → Fault Lines table → Drills 1–3
Targeted Study1–2 hrsExecutive Summary → Walkthrough → Section 3 (Fault Lines) → Section 4 (Failure Modes) → Deep Dives 1, 2 and 4
Deep Dive3+ hrsEverything, including Section 11 (Principal Lens) and Appendix A (lease-based delivery)
What is a Distributed Message Queue? — Why interviewers pick this topic

A distributed message queue accepts messages from producers, stores them durably across machines, and hands them to consumers — later, possibly much later, possibly more than once. It decouples when work is requested from when it is done, and who requests it from who does it.

That decoupling is the whole value and the whole danger. The producer has been told "accepted" and has moved on. From that moment, the queue owns the message: if a broker dies, the message must survive; if the consumer crashes mid-processing, the message must come back; if the message can never be processed, someone must find out.

Before vs After — the "order confirmation emails" scenario:

Without a designed delivery contract:
t=0:        Checkout publishes "order_placed" to the queue. Broker acks after
            writing to the leader's page cache only.
t=+40ms:    Leader host loses power. Follower had not fetched the message yet.
t=+12s:     Follower promoted. The message never existed. Producer already returned 200.
t=+3 days:  Customer emails support: "no confirmation". Nobody can find the message.
            Meanwhile, a consumer that takes 45s per message has a 30s visibility
            timeout — every slow message is processed twice. 3% of customers get
            duplicate emails. The team "fixes" it by raising the timeout to 12h,
            so the next consumer crash strands 8,000 messages for 12 hours.

With an explicit contract:
t=0:        Producer sends with idempotence key. Broker acks only after the message
            is on 2 of 3 replicas in 2 AZs (p99 ~8ms).
t=+40ms:    Leader dies. Follower already has the message. Promoted in ~5s.
t=+5s:      Producer retries in-flight batch. Broker dedups by (producer_id, seq).
t=+6s:      Consumer leases message for 60s, heartbeats every 20s while working.
t=+51s:     Consumer finishes, acks. No duplicate. Lease never expired.
t=+1 day:   One malformed message failed 5 times → DLQ. Owning team paged
            by dlq.oldest_message_age > 1h. Fixed and redriven the same day.

Why interviewers reach for this question: A queue is the smallest system that forces every hard distributed-systems decision at once — durability vs latency at ack time, ordering vs parallelism, at-least-once vs "exactly-once", push vs pull, and the ownership of failure. A candidate who says "we'll use Kafka" has picked a component. The interviewer wants to know whether they can define the contract that component must honor, and what happens to the 0.01% of messages that break it.

Mechanics Refresher: Storage and Delivery Models
ModelHow It WorksProsCons
Per-message state queue (SQS, RabbitMQ classic)Each message has its own state: visible, in-flight (leased), deleted. Consumers receive, then deleteIndependent per-message ack; consumers scale freely; natural DLQPer-message bookkeeping is expensive; no replay after delete; ordering is weak or costly
Partitioned log (Kafka, Pulsar storage, Kinesis)Append-only log per partition; consumers track an offset; messages retained for a windowVery high throughput (sequential I/O); replay; many independent consumer groupsParallelism capped at partition count; one stuck message blocks its partition; no per-message ack
Log + lease layer (Pulsar shared subscriptions, consumer proxies over Kafka)Log for storage; a dispatcher tracks per-message acks on top of offsetsReplay and per-message ack; parallelism beyond partitionsMore moving parts; ack-state must itself be durable
Database-as-queue (SELECT … FOR UPDATE SKIP LOCKED)Rows are jobs; workers lock and deleteTransactional with business data; zero new infraTable bloat and vacuum pressure past ~1–5K jobs/s; polling load
In-memory broker (Redis lists/streams)Messages in RAM, optional AOFSub-millisecond; simpleDurability is a config flag, not a guarantee; memory-bound backlog

For most production systems: a partitioned, replicated log for storage, with a lease-based delivery layer on top when consumers need per-message acknowledgement, retries and DLQs. The storage engine is not the interview — the delivery contract is.


Executive Summary

If you only read one section, read this. Everything in this case study flows from the contrast below.

What This Interview Actually Tests#

A message queue is not a storage question. Everyone can append bytes to a file.

It is a contract question: what exactly has the producer been promised when it gets an ack, what exactly has the consumer been promised when it gets a message, and who is accountable for the messages that fall between those two promises. It tests:

  • Whether you define durability at ack time in terms of replicas and failure domains, not adjectives
  • Whether you know that ordering and parallelism are the same knob turned in opposite directions
  • Whether you treat redelivery as normal operation and push idempotency to the consumer, by contract
  • Whether you design for the backlog — the moment consumers fall behind is when the queue earns its keep or loses data
  • Whether the DLQ has an owner, an alert and a redrive path, or is a place messages go to die

The key insight: A queue does not remove failure; it moves failure in time. Every message you accept is a promise you'll keep later, under conditions you can't see now. Staff candidates design the conditions under which that promise is kept — ack quorum, lease duration, retry budget, retention vs maximum tolerable lag — and name who gets paged when it isn't.

The L5 vs L6 Contrast — Start Here#

BehaviorSenior (L5)Staff (L6)Principal (L7)
First moveDraws producer → broker cluster → consumer; picks KafkaAsks "Is this a work queue, an event log, or fan-out? Do consumers need per-message ack, replay, or both?"Asks "How many queueing systems does the org already run, and which team will own this one's pager for five years?"
Durability"Replication factor 3""The ack means: written to 2 of 3 replicas in distinct AZs. RF=3 with ack-on-leader is still a single copy for ~10–50ms"Sets the org durability tiers (ephemeral / standard / critical) with a cost per tier, so teams stop defaulting to the most expensive one
Delivery semantics"Exactly-once with Kafka transactions""At-least-once delivery, idempotent effect at the consumer keyed by message ID; transactions only for read-process-write within the log"Makes the consumer idempotency contract a platform rule enforced by a client library and a conformance test, not a wiki page
Ordering"Use one partition to keep order""Order per entity key only; the key is the unit of ordering and of blocking. Poison messages park the key, not the partition"Recognizes global ordering requests as a product requirement in disguise and pushes back with the cost: one partition = ~10–50 MB/s ceiling, forever
Failure"Add retries and a DLQ""Retry budget of 5 with exponential backoff, then DLQ; DLQ has an owner, a 1h age alert, and a redrive tool that preserves keys"Defines the org's backlog posture: max tolerable lag per tier vs retention, quarterly consumer-outage game days, error budget on end-to-end latency
Scale"Add more partitions and brokers""Partition count is a 2-year decision for keyed topics — repartitioning breaks key ordering. Size for 3× peak, cap partition size at ~25 GB"Plans the 3-year path from shared cluster to cells per business domain, with a tenancy model and chargeback
Why "delivery semantics" separates levels

L5: "We'll enable exactly-once semantics so consumers never see duplicates." This is a real feature in several brokers, and the candidate is not wrong that it exists. But it covers a narrow case: reading from the log, transforming, and writing back to the log in one transaction. The moment a consumer sends an email, charges a card or writes to an external database, the broker's transaction boundary ends and duplicates are possible again.

L6: "The queue guarantees at-least-once delivery. Duplicates happen on consumer crash after the side effect but before the ack, on lease expiry, and on producer retry after a lost ack. Every message carries a stable message_id; consumers record it in the same transaction as their side effect — a unique index on processed_messages(message_id) — so a redelivery is a no-op. That gives us exactly-once effect where it matters, without pretending the network is reliable."

L7: "The duplicate isn't the queue team's bug or the consumer team's bug — it's an unowned contract. I'd ship the dedup as part of the consumer SDK with a pluggable store, make 'idempotent handler' a checkbox in the service readiness review, and run a monthly chaos test that redelivers 1% of messages in staging. Teams that fail it don't get access to critical-tier topics."

Why "ordering" separates levels

L5: "Orders for the same customer must be processed in order, so we'll put them in one partition." Reasonable instinct. But it says nothing about what happens when message #4 for customer 812 fails: either everything behind it on that partition waits (head-of-line blocking) or the consumer skips it and the order guarantee is quietly broken.

L6: "Ordering is per key — account_id. The partition is just a container for many keys. When a message for one key fails after retries, I park that key: subsequent messages for account 812 go to a hold area, but the other ~20,000 accounts on that partition keep flowing. I'd rather delay one customer by an hour than every customer on the partition."

L7: "Before I commit to ordering at all, I'll ask the product owner which business invariant actually needs it. Nine times out of ten, the consumer can tolerate reordering if events carry a version number and the handler ignores stale versions. Ordering is the most expensive guarantee a queue makes; I want it to be a deliberate purchase with a named buyer."

Why "failure" separates levels

L5: "Failed messages go to a dead-letter queue." The DLQ exists. Nobody watches it. Retention is 14 days. On day 15, messages start silently expiring, and the first person to notice is a customer.

L6: "The DLQ is a production queue with an owner: the consuming team. Alert on dlq.depth > 0 for critical topics and dlq.oldest_message_age > 1h for everything else. The redrive tool replays into the original queue preserving the key and message ID, with a rate limit so we don't stampede the fixed consumer. And DLQ retention is longer than primary retention — 14 days vs 4 — because messages land there precisely when humans are slowest to react."

L7: "Across the org, I'd track one number: messages that expired unprocessed, per month, per team. It should be zero for critical tiers. When it isn't, it shows up in the same review as SLO misses, because a lost message is an outage that nobody saw."

The Staff Positions#

PositionRationale
Ack after quorum replication across failure domains, not after leader writeThe ack is the moment the queue takes ownership. Ack-on-leader turns any leader crash in the replication window (~10–50ms) into silent loss
At-least-once delivery + idempotent consumers, by contractExactly-once delivery across a side effect is not achievable; exactly-once effect is, and it lives in the consumer's transaction
Order per key, never per queueGlobal order caps throughput at one partition and turns every poison message into a full outage
Pull with long-polling for work queues; consumers own their concurrencyThe consumer knows its own capacity; push-based brokers have to guess and overrun slow consumers
Leases (visibility timeouts) with heartbeats, not fixed long timeoutsA short lease (30–60s) plus heartbeat extension gives fast redelivery on crash without duplicating slow work
Bounded retries, then DLQ with an owner and an age alertInfinite retries convert one poison message into permanent load; an unowned DLQ converts it into silent loss
Retention ≥ 3× the longest plausible consumer outageBacklog that outlives retention is data loss. If a consumer can be down for a weekend (~60h), 4-day retention is too short

The Three Intents#

Three intents produce three different systems. Name them, then commit.

IntentConstraintStrategyFailure ModeCorrectness Bar
Work queue (async jobs between services: send email, resize image, charge card)Each message processed successfully once; consumers scale with load; slow/poison messages isolatedPer-message leases with visibility timeout, bounded retries, DLQ, optional per-key orderingDuplicate side effects on lease expiry; poison messages; backlog outliving retentionZero acked-then-lost messages; duplicate effect rate 0 with idempotent handlers; DLQ age < 1h
Event log (domain events consumed by many teams: order_placed, CDC streams)Replay; many independent consumer groups; per-entity orderPartitioned replicated log, offsets per group, long retention or compactionHot partitions; consumer lag; schema breakage across teamsPer-key order preserved; replayable for the retention window; schema compatibility enforced
Ephemeral fan-out (presence, live dashboards, cache invalidation)Latency < 50ms; many subscribers; old messages worthlessIn-memory pub/sub, no persistence or short TTL, drop on slow subscriberSlow subscriber backs up the broker; message stormsFreshness over completeness; loss acceptable and documented

🎯 Staff Move: "I'll design the work-queue intent — durable async jobs between services, at-least-once, per-message ack, retries and DLQ — because that's where delivery semantics, leases and poison handling all bite. I'll build it on a partitioned replicated log so we keep replay and can serve the event-log intent later with the same storage. If you meant fan-out for live updates, that's a different, much lighter system and I'd drop durability entirely."

The Five Fault Lines#

#Fault LineThe Tension
1Ack Durability vs Producer LatencyAck after leader write (~1–2ms) or after cross-AZ quorum (~5–15ms)? Who loses data when the leader dies in between?
2Ordering vs ParallelismEvery ordering guarantee is a cap on concurrency and a head-of-line blocking risk. How much order is the business actually buying?
3Redelivery Speed vs Duplicate WorkShort leases redeliver fast after a crash but duplicate slow work; long leases strand messages when consumers die
4Retry Forever vs Dead-LetterKeep retrying (no human needed, but poison loops) or give up after N (bounded load, but now someone must own the DLQ)?
5Shared Multi-Tenant Cluster vs Dedicated QueuesEfficiency and one operating team, or isolation from the noisy neighbor that eats your partition's disk and I/O?

In the Wild: Real Production Systems#

Why this section belongs here: Naming how real systems settled these fault lines shows you've studied operations, not just APIs.

Amazon SQS — The Lease Model as a Public Contract#

SQS standard queues are documented as at-least-once with best-effort ordering; a received message becomes invisible for a visibility timeout (default 30 seconds, configurable up to 12 hours) and reappears if not deleted. A redrive policy moves a message to a dead-letter queue after maxReceiveCount receives. FIFO queues add ordering within a message group ID and deduplicate by a deduplication ID within a 5-minute window. Retention defaults to 4 days, maximum 14.

Staff insight: SQS makes every fault line in this case study a visible parameter — lease length, receive count, group key, dedup window, retention. In an interview, naming those knobs and their defaults signals you've run one; explaining why the dedup window is only 5 minutes (dedup state is expensive, and producer retries happen within seconds) signals Staff.

Uber — A Consumer Proxy in Front of Kafka#

Uber has written publicly about a Kafka consumer proxy that pulls from Kafka and pushes to service endpoints over gRPC. The proxy tracks per-message acknowledgement, retries failed messages through retry queues, routes exhausted ones to a DLQ, and lets a service process a single partition with many parallel workers — decoupling consumer parallelism from partition count.

Staff insight: This is the "log + lease layer" model. The log gives durability and replay; the proxy gives work-queue semantics and isolates poison messages without stalling the partition. It's the shape this case study's design converges on, and it's the answer to "but Kafka can't do per-message ack."

Slack — Durable Ingestion in Front of a Fast Job Queue#

Slack has described re-architecting its asynchronous job queue after a Redis-based queue hit memory exhaustion during a backlog and lost the ability to accept jobs. The redesign put Kafka in front as a durable buffer, with a relay feeding jobs into Redis for execution, so a backlog grows on disk rather than in RAM.

Staff insight: The failure was not "Redis is bad" — it was the backlog had nowhere to go. Queues fail at exactly the moment they're needed most: when consumers are slow. A Staff design states where the backlog lives at 10× normal depth and how long it can grow before something breaks.

What Interviewers Probe#

After You Say...They Will Ask...(What They're Evaluating)
"Replication factor 3""When exactly do you ack the producer? What if the leader dies 20ms after the ack?"Ack semantics vs replica count
"Consumers process exactly once""The consumer charged the card, then crashed before acking. What happens?"Exactly-once effect vs delivery
"We keep messages in order""Message 4 of 10 for this customer fails forever. What happens to 5–10? To other customers?"Head-of-line blocking, poison handling
"Visibility timeout of 30 seconds""Some jobs take 3 minutes. Some consumers die. Pick one timeout."Leases + heartbeats
"Failed messages go to a DLQ""Who looks at it? What happens on day 15?"DLQ ownership, retention
"We'll autoscale consumers""The consumers write to a database that's already at 80% CPU. Now what?"Backpressure, downstream capacity
"We'll add partitions when we need more throughput""The topic is keyed by account. What happens to ordering when you go from 32 to 64?"Repartitioning as a one-way door

System Architecture Overview#

Diagram: System Architecture Overview

Quick-Reference: The 30-Second Cheat Sheet#

QuestionStaff Answer
What does the ack mean?Message is on 2 of 3 replicas in different AZs; survives any single AZ loss
Delivery semantics?At-least-once; exactly-once effect via consumer-side dedup on message_id
Ordering?Per key, opt-in; poison messages park the key, not the partition
Push or pull?Pull with long-poll (20s) for work queues; consumer owns concurrency
Crash mid-processing?Lease expires (60s) → redelivered; heartbeats extend leases for long jobs
Poison message?5 attempts with backoff (10s → 1m → 10m), then DLQ with owner + age alert
Backlog?Retention ≥ 3× worst consumer outage; alert on oldest_unacked_age, not depth
Scale?Partitions sized for 3× peak; ~5–10 MB/s per partition budget; never repartition keyed topics casually
Multi-tenant?Per-tenant produce/consume quotas at the front end; cells for blast radius

Key Numbers Worth Memorizing#

NumberValueContext
Leader-only ack latency~1–2 msPage-cache write, same host; no durability across host loss
Cross-AZ quorum ack latency~5–15 ms p99One RTT to the fastest follower (~1–2ms) + batching + fsync policy
Per-partition throughput budget~5–10 MB/s planned, 30–50 MB/s ceilingPlan at a fraction of the ceiling so a single consumer can keep up
Leader failover~2–10 sDetection (session timeout) + election + client metadata refresh
SQS visibility timeout30 s default, 12 h maxThe canonical lease parameter
SQS retention4 days default, 14 days maxSets the outer bound on tolerable consumer outage
SQS long-poll waitup to 20 sCuts empty receives (and their cost) by ~90%+ on idle queues
SQS FIFO dedup window5 minDedup state is bounded; producer retries happen within seconds
Kafka min.insync.replicas for critical data2 (with RF=3)Writes fail rather than ack with a single copy
Typical message size1–10 KBPayloads > 256 KB–1 MB belong in blob storage with a pointer
Retry schedule10 s, 1 m, 10 m, 1 h, then DLQBounded total ~1.2 h; covers most transient downstream outages
DLQ age alert1 h (critical: any message)DLQ is a production queue, not an archive
Database-as-queue ceiling~1–5K jobs/sBefore vacuum/bloat and lock contention dominate on Postgres
Replication bandwidth(RF − 1) × produce rate100 MB/s in at RF=3 = 200 MB/s cross-AZ, billed per GB on most clouds

Interview Walkthrough

The walkthrough below is a 45-minute script. The words in italics are meant to be said out loud.

Phase 1: Requirements & Framing (2–3 minutes)#

Open by splitting the intents, then commit.

"Before I draw anything — 'message queue' covers three different systems. A work queue where each job is done once and acked individually; an event log that many teams replay; and ephemeral fan-out where old messages are worthless. I'll design the work queue as an internal platform: durable, at-least-once, per-message ack, retries, DLQ, and opt-in per-key ordering. I'll put it on a partitioned log so replay is available later."

Then pin down the numbers that change the design:

QuestionAssumed AnswerWhy It Matters
Peak produce rate?200K msg/s across all tenants, 2 KB average → ~400 MB/sSets partition count and replication bandwidth
Message size distribution?p50 1 KB, p99 64 KB, hard cap 1 MBLarge payloads go to blob storage; changes batching
Processing time per message?p50 200 ms, p99 30 s, some jobs 10 minDecides lease length and heartbeat design
Ordering needed?Some tenants need per-key order (account events); most don'tOrdering is opt-in per queue
Loss tolerance?Zero after ack for critical tierAck after cross-AZ quorum
Duplicate tolerance?Consumers must tolerate duplicatesAt-least-once, idempotency contract
Longest consumer outage to survive?72 h (long weekend)Retention ≥ 7 days on critical tier
Tenancy?~300 internal teams, ~5,000 queuesQuotas, isolation, cells

"I'll treat end-to-end latency — produce to consume — as p99 under 100ms when consumers are healthy, and I won't optimize below that; durability and isolation matter more here."

Phase 2: Core Entities & API (1–2 minutes)#

EntityKey FieldsNotes
Queuequeue_id, tenant, tier, partitions, ordering_mode, retention, max_attempts, dlq_idConfig is versioned; changes are rolled out like code
Messagemessage_id (producer-supplied or broker-assigned), key, payload, headers, enqueued_at, attemptmessage_id is the idempotency handle for consumers
Leasemessage_id, consumer_id, lease_expires_at, receipt handleReceipt handle is per-delivery, so a stale consumer can't ack a re-leased message
Partitionpartition_id, leader, replicas, high-watermarkInternal; not exposed in the API
POST /queues/{q}/messages                       # produce (batch up to 500 / 1 MB)
  body: [{ message_id?, key?, payload, headers, delay_seconds? }]
  → 200 { results: [{ message_id, status: "accepted" | "duplicate" }] }

POST /queues/{q}/receive                        # long-poll
  body: { max_messages: 10, wait_seconds: 20, lease_seconds: 60 }
  → 200 [{ message_id, receipt, payload, attempt, enqueued_at }]

POST /queues/{q}/ack        { receipts: [...] }
POST /queues/{q}/nack       { receipt, retry_after_seconds? }
POST /queues/{q}/extend     { receipt, lease_seconds }   # heartbeat
POST /queues/{q}/redrive    { from: dlq, rate_per_sec, filter? }   # operator API

"Two details matter: the receipt is per-delivery, not per-message — a consumer whose lease expired can't ack a message someone else now holds. And produce returns 'duplicate' rather than an error when a producer retries a message ID we've already accepted, so the producer's retry loop is trivially safe."

Phase 3: High-Level Architecture (≤5 minutes)#

Staff candidates spend under five minutes here. Draw the tiers, state what each one owns, and move on.

Diagram: Phase 3: High-Level Architecture (≤5 minutes)

"Front-end is stateless: auth, per-tenant quotas, and routing by hash(key) mod partitions. Storage is a replicated log — I'll lean on Kafka-style internals and not redesign segment files. The interesting tier is the dispatcher: it turns a log into a work queue by tracking per-message leases on top of the offset, so one slow message doesn't stall a partition and consumers can scale past partition count. Metadata — partition maps, leaders, queue config — lives in a small Raft group, like etcd."

Phase 4: Transition to Depth (1 minute)#

The sentence that steers the interviewer toward your strongest ground:

"The boxes are standard. The design lives in three places: what the producer ack means when a leader dies, what happens between a consumer receiving a message and acking it, and what happens to messages that can never succeed. I'd like to go through those in that order, then cover multi-tenancy. Is there one you'd rather start with?"

Phase 5: Deep Dives (25–30 minutes)#

Deep dive 1 — The ack path (8 min).

Diagram: Phase 5: Deep Dives (25–30 minutes)

"The ack goes back when the message is on the leader and at least one follower in another AZ — min_in_sync = 2 of RF 3. That costs one cross-AZ round trip, ~1–2ms, plus batching; p99 lands around 5–15ms. If the in-sync set drops below 2, I reject produces rather than ack with one copy: producers see errors, which is honest, instead of success, which is a lie. Producer idempotence — producer ID plus sequence number per partition — means a retry after a lost ack is deduplicated by the leader."

"On fsync: I don't fsync per message. Two copies in two AZs' page caches is stronger than one copy fsynced, because correlated failure across AZs is much rarer than a single host losing power. For the critical tier, I'd add periodic fsync — every 100–500ms — to bound loss under a full-region power event."

Deep dive 2 — The lease (10 min).

Diagram: Phase 5: Deep Dives (25–30 minutes)

"A receive grants a lease of 60s by default. Long-running consumers heartbeat via extend every lease/3 — 20s — so a 10-minute job never gets redelivered while its worker is alive, but a crashed worker's message comes back within 60s. The receipt handle includes a lease generation; an ack with a stale generation is rejected and counted in ack.stale_receipt_rate, which is my early warning that leases are too short."

"The dispatcher tracks leases per partition. Its state — which offsets are acked, which are leased until when — goes to a compacted ack log, so a dispatcher failover rebuilds in seconds. The committed offset for the partition is the lowest un-acked message; everything below it is done. If one message is stuck at attempt 3 while 50,000 after it are acked, the gap is tracked as a sparse set, not by blocking."

Deep dive 3 — Poison messages and DLQ (6 min).

"Each nack or lease expiry increments attempt. Retries are delayed, not immediate: 10s, 1m, 10m, 1h — implemented as delay tiers rather than per-message timers, so the delay mechanism is just more queues. After attempt 5, the message moves to the DLQ with its full history: error strings, consumer IDs, timestamps. The DLQ belongs to the consuming team, alerts on age, and has a rate-limited redrive. For ordered queues, a dead-lettered message parks its key: later messages for that key are held until the operator redrives or explicitly skips — because skipping silently is how you corrupt an account's state."

Deep dive 4 — Multi-tenancy (4 min).

"Per-tenant quotas at the front end — bytes/s and requests/s for produce and consume — enforced with local token buckets synced every second. Storage quotas per queue cap the backlog: when a tenant hits it, they get 429s, not the cluster. Critical-tier tenants get placement in a separate cell — separate brokers, separate metadata group — so a noisy best-effort tenant can't fill their disks."

Phase 6: Wrap-Up (2–3 minutes)#

"To summarize the contract: an ack means two AZs have it; a delivery means you hold a lease and must be idempotent; a failure means five attempts over about an hour, then a DLQ your team owns. What I didn't cover: cross-region replication — I'd do async mirroring with consumer-offset translation and accept a seconds-scale RPO — and schema governance for payloads. What I'd build next: per-key parking for ordered queues and a self-service redrive UI, because the DLQ is where the human cost of this system lives."

Common Timing Mistakes#

MistakeTime LostFix
Designing the log segment format and index files10+ minSay "standard segmented log, see Kafka internals" and move on
Debating Kafka vs RabbitMQ vs SQS before defining requirements5 minDefine the contract first; the product choice follows
Drawing ZooKeeper election in detail5 min"Metadata in a Raft group; leader election is a solved problem here"
Never reaching consumer failure or DLQWhole interviewSteer explicitly in Phase 4 — that's where the Staff-level signal is
Sizing partitions to three decimal places3 minOne line: "~400 MB/s, ~10 MB/s/partition, ~120 partitions at 3× headroom"

1. The Staff Lens#

1.1 Why This Problem Exists in Staff Interviews#

Almost every large design — notifications, payments, feeds, ad aggregation, search indexing — has a queue in the middle. Interviewers use the queue question to test whether a candidate understands the guarantees they're assuming every time they draw that arrow. A Senior engineer has used queues heavily; a Staff engineer has been paged because one silently dropped messages, redelivered 40,000 of them in a storm, or let a single malformed payload stall a partition for six hours. The question separates users of queues from owners of them.

It's also an ownership question in disguise. A queue sits exactly on a team boundary: the producer team, the queue team and the consumer team each believe someone else is responsible for the message that fails. The Staff answer draws those boundaries explicitly.

1.2 The L5 vs L6 Contrast — Visual#

Diagram: 1.2 The L5 vs L6 Contrast — Visual

The Senior path is a sequence of component choices. The Staff path is a sequence of promises, each of which has a cost and a failure mode.

1.3 The Staff Question That Cuts Through Everything#

"When a consumer receives this message, what has it been promised — and what has it promised back?"

The answer forces every decision into the open: whether the message can arrive twice (yes), whether it can arrive out of order (depends on the key), how long the consumer has before someone else gets it (the lease), and what the consumer must guarantee (idempotent effect, ack after side effect, heartbeat if slow). A candidate who can answer that question crisply has designed the queue; the rest is capacity planning.

🎯 Staff Move: "Let me write down the consumer contract before the architecture: you may see a message more than once, you have 60 seconds unless you heartbeat, you must ack after your side effect commits, and after five failures it goes to your DLQ. Everything I build is there to make that contract true."


2. Problem Framing & Intent#

2.1 The Three Intents — Explained#

Work queue. The unit of value is a completed job. Producers don't care which consumer does the work or exactly when; they care that it's done once, eventually, and that failures surface. The design centers on per-message acknowledgement, leases, retries and the DLQ. Ordering is usually unnecessary and, where it's needed, is per entity. Messages have no value after they're processed — retention exists only to survive consumer outages.

Event log. The unit of value is the history. Many consumer groups — search indexing, analytics, fraud, notifications — read the same events independently and at their own pace; a new consumer may need to replay last week. The design centers on partitioning, per-key order, retention or compaction, and schema contracts. Per-message ack is a non-goal; an offset per partition per group is enough. This is the Kafka sweet spot and the backbone of stream processing.

Ephemeral fan-out. The unit of value is now. A typing indicator from three seconds ago is noise. The design centers on low latency and dropping data for slow subscribers instead of buffering it. Persistence is a liability: it adds latency and creates backlogs nobody wants. See Real-time WebSockets for where this lives in practice.

DimensionWork QueueEvent LogEphemeral Fan-out
Ack granularityPer messagePer partition offsetNone
ReplayNot requiredCore featureNever
Consumer scalingUnbounded≤ partition count per groupPer subscriber
Slow consumerBacklog grows; alert on ageLag grows; alert on lagDrop messages for that subscriber
Retention driverMax consumer outageReplay window, new consumersSeconds or zero
Failure ownerConsuming team (DLQ)Each consumer groupNobody — loss is accepted

2.2 When NOT to Use a Distributed Message Queue#

SituationBetter ChoiceWhy
Caller needs the result to respond to the userSynchronous RPC with timeout and retryA queue adds latency and a second failure path, and the caller waits anyway
< 1K jobs/s and the job must commit with business dataSKIP LOCKED table in PostgreSQLTransactional enqueue, zero new infrastructure; the outbox pattern is a table already
Scheduled execution at a precise future time (reminders in 30 days)A scheduler with a time index — see Task SchedulingQueues are bad at "deliver at 09:00 on March 3"; delay caps are usually 15 min–12 h
Long-running, multi-step business processes with compensationWorkflow engine (durable execution) — see Long-Running ProcessesChaining queues re-invents a state machine without visibility
Live updates where stale data is worthlessIn-memory pub/sub or direct WebSocket fan-outDurability adds latency and backlog with no value
"Decoupling" two services owned by the same team that deploy togetherA function callA queue at a non-boundary is operational cost for no organizational benefit

🎯 Staff Move: "Before adding a queue, I'd ask whether anyone needs the work done later rather than now. If the caller blocks on the result, a queue is just a slower, less observable RPC."

2.3 What the Interviewer Leaves Underspecified#

Unstated AssumptionWhy It's Deliberately VagueWhat to Say
Delivery semanticsTo see if you say "exactly-once" without qualification"At-least-once delivery; exactly-once effect is the consumer's job, and I'll give them the tools."
Ordering scopeTo see if you default to global order"Ordering per key, opt-in. Most queues won't need it."
Processing timeTo see if you notice leases depend on it"What's the p99 handler time? Lease length follows from it."
Consumer outage durationTo see if you link retention to operations"How long can a consumer team be down before we lose data? That sets retention."
Who owns failuresTo see if you think organizationally"The consuming team owns its DLQ; the platform owns the tooling and the alert."
Message sizeTo see if you push large payloads elsewhere"Anything over 256 KB goes to blob storage; the message carries a pointer."
Number of tenantsTo see if you consider isolation"Shared cluster or per-team? That decides quotas versus cells."

2.4 Precise Terminology#

TermPrecise MeaningCommon Misuse
Ack (producer)The broker accepted responsibility; message survives the failures the durability tier covers"Message was written somewhere"
Ack (consumer)The consumer's side effect is committed; the message may be deletedAcking on receipt, before processing
At-least-onceEvery acked message is delivered ≥1 time; duplicates possibleUsed as if duplicates were rare edge cases — they're routine on every crash
Exactly-once effectDuplicates are delivered but have no additional effect, via idempotent consumers"Exactly-once delivery," which no system provides across external side effects
Lease / visibility timeoutTime-bounded exclusive claim on a message"Lock" — a lease expires without the holder's cooperation
Head-of-line blockingA stuck message prevents later messages in the same ordering scope from being processedBlamed on "slow consumers" when the cause is one poison message
Poison messageA message that fails deterministically regardless of retriesConfused with transient failure — retrying a poison message is pure waste
Consumer lagMessages (or time) between the latest produced and the latest committedMeasured in messages only; age of oldest unacked message is what matters
BackpressureSignal from a slower stage that limits the rate of a faster stage"Autoscaling," which increases rate rather than limiting it
In-sync replicaA follower within the replication-lag bound that can be promoted without lossAny replica that exists

3. The Five Fault Lines#

Each fault line below is a place where two competent engineers can disagree. The Staff contribution is not the "right" answer — it's naming who pays for each option and committing to a default with a stated reason to deviate.

3.1 Fault Line 1: Ack Durability vs Producer Latency#

The producer ack is the moment ownership transfers. Before it, losing the message is the producer's problem (it will retry). After it, losing the message is the queue's fault — and nobody will notice until a customer does.

Ack PolicyProducer p99What SurvivesWhat BreaksWho Pays
Fire-and-forget (no ack)< 1 msNothing guaranteedAny network blip loses data silentlyDownstream teams, who debug "missing events" for weeks
Leader write (page cache)~1–2 msProcess crash on leader (page cache persists)Leader host loss in the replication window: acked data gone; unclean election makes it worseCustomers whose messages vanish; the on-call who can't prove what was lost
Leader fsync~2–10 ms (SSD), batchedLeader host power lossLeader disk/AZ loss — single copySame, less often; latency paid by every producer
Quorum: 2 of 3, cross-AZ~5–15 msAny single host or AZ lossSimultaneous loss of two AZs before flush (very rare)Producers pay ~5–10 ms; cluster pays (RF−1)× cross-AZ bandwidth
All replicasTail of the slowest replica, 20–100 ms+Same as quorumOne slow replica stalls every producerEvery producer, every time a disk hiccups

The Staff default: quorum ack, min_in_sync = 2 of RF 3, replicas spread across three AZs, with produce rejected when fewer than two replicas are in sync. Periodic fsync every 100–500 ms for the critical tier.

When to deviate: metrics, logs and clickstream where 0.01% loss is cheaper than 5 ms per call and 2× bandwidth — offer an ephemeral tier with leader ack and RF=2. Make it a named tier with a documented loss bound, never a per-producer config flag nobody reviews.

🎯 Staff Move: "Replication factor is how many copies eventually exist. The ack policy is how many exist when I tell the producer 'you can forget this.' I'll specify the second — two copies in two AZs — and I'll reject writes rather than ack with one."

Full reasoning: why rejecting writes beats acking with one copy

When two of three replicas fall out of sync (a slow disk plus a deploy, say), the leader can either keep accepting writes with one copy or refuse them. Accepting looks like availability, but it converts a visible outage (producers get errors and retry, or buffer locally) into an invisible one (data that will be lost if this single host fails in the next minutes). Producers have retry loops and local buffers precisely for transient unavailability; they have no mechanism at all for "the ack lied."

The cost: during an under-replicated window, producer error rates climb. That needs an alert (partition.under_min_isr > 0 for 2 min pages the storage on-call) and producers sized to buffer ~60 seconds of traffic. Kafka's acks=all with min.insync.replicas=2 implements exactly this policy; the design choice transfers to any replicated log.

3.2 Fault Line 2: Ordering vs Parallelism#

Ordering and parallelism are the same knob. Any message that must wait for another is a message that can't be processed concurrently. Every ordering guarantee also defines a blast radius for poison messages: everything in the same ordering scope waits behind it.

Ordering ScopeMax ParallelismPoison Blast RadiusWho Pays
Global (one queue)1 consumerEntire queue stopsEvery consumer and every downstream user; throughput capped at ~10–50 MB/s forever
Per partition= partition countEvery key on that partition (often 10K+ entities)Unrelated customers who happen to hash to the same partition
Per key, partition-blocking= partition countSame as per partition — the common accidental designSame; looks per-key on paper
Per key, key-parkingUnbounded (one in-flight per key)One keyThe one customer whose message is broken
NoneUnboundedOne messageConsumers that needed order and must now use versioning

The Staff default: no ordering unless requested; when requested, per key with one in-flight message per key and key-parking on poison. The dispatcher tracks in-flight keys per partition; a key with a dead-lettered message is held until redrive or explicit skip.

When to deviate: if consumers can tolerate reordering by carrying a version or timestamp (last-writer-wins on entity_version), drop ordering entirely — it's cheaper and removes head-of-line blocking. If a tenant genuinely needs total order (a ledger posting stream), give them a dedicated single-partition queue and document the throughput ceiling in their SLO.

Diagram: 3.2 Fault Line 2: Ordering vs Parallelism

🎯 Staff Move: "I want ordering to be a purchase, not a default. If you need it, you get it per key, one in flight per key, and a poisoned key parks itself — I'd rather delay one account than every account that hashes next to it."

3.3 Fault Line 3: Redelivery Speed vs Duplicate Work#

The lease is a bet on how long processing takes. Bet too short and slow-but-healthy consumers get their messages yanked and redelivered to someone else — duplicate work, duplicate side effects, and sometimes a feedback loop. Bet too long and a crashed consumer's messages sit invisible until the lease expires.

Lease StrategyCrash Recovery TimeDuplicate RateWhat BreaksWho Pays
Short fixed (30 s)≤ 30 sHigh for any job with p99 > 30 sSlow jobs run 2–3× in parallel; downstream load multipliesDownstream services (load), customers (duplicate emails/charges if not idempotent)
Long fixed (12 h)≤ 12 hNear zeroCrashed consumer strands messages for hours; backlog invisibleUsers waiting for work stranded in a dead worker
Short + heartbeat extension≤ lease (60 s)Near zero for live consumersConsumer that hangs while its heartbeat thread lives holds messages foreverNeed a max total lease (e.g., 1 h) and handler-level watchdogs
Adaptive (lease = k × p99 handler time)VariesLowComplex; p99 shifts during incidents exactly when you need stabilityPlatform team maintaining the model

The Staff default: 60 s lease, heartbeat every lease/3, maximum total lease of 1 h (configurable per queue up to 12 h), and stale-receipt acks rejected. A worker that can't finish in the max lease should be a workflow, not a message.

When to deviate: for very short jobs (p99 < 1 s) at high volume, a 10–15 s lease without heartbeats reduces dispatcher load. For batch jobs that legitimately take hours, move them to a job scheduler with checkpointing.

🎯 Staff Move: "The lease should reflect how fast I want to notice a dead consumer, not how long the slowest job takes. Heartbeats decouple the two: 60 seconds to detect a crash, and a job can run as long as it keeps proving it's alive — up to an hour."

3.4 Fault Line 4: Retry Forever vs Dead-Letter#

Every failed message is either retried or set aside. Retrying requires no human but can loop forever; setting aside bounds load but requires an owner. The question is never "should we have a DLQ" — it's who looks at it and how fast.

PolicyLoad From PoisonHuman RequiredWhat BreaksWho Pays
Immediate retry, unbounded100% of consumer capacity per poison message, foreverNoOne bad message burns a worker permanently; storms during downstream outagesConsumer fleet; downstream dependency during its recovery
Backoff retry, unboundedSmall but permanentNoPoison messages accumulate; retention eventually expires them silentlyCustomers whose messages expire; nobody sees it
Bounded retry → DLQ, unownedBoundedNobody, in practiceDLQ fills; messages age out on day 14Same customers, with a false sense of safety
Bounded retry → DLQ, owned + alertedBoundedYes, the consuming teamToil if the failure rate is highConsuming team's on-call — correctly, because they own the fix
Bounded retry → drop + logBoundedNoData loss by designAcceptable only for ephemeral-tier messages with product sign-off

The Staff default: five attempts with backoff (10 s, 1 m, 10 m, 1 h), classify errors where possible (a 400 Bad Request from the handler goes straight to DLQ — retrying a deterministic failure is waste), DLQ owned by the consuming team, alert on age, retention 14 days in the DLQ vs 4–7 in the primary.

When to deviate: if the downstream dependency has a known long outage window (a partner API with maintenance windows), extend the retry budget to cover it rather than dead-lettering thousands of good messages; the tell is a DLQ where >90% of messages redrive successfully without code changes.

🎯 Staff Move: "A DLQ is a queue with a pager attached, or it's a slower way of deleting data. I'll make the consuming team its owner, alert on the age of the oldest message, and give them a redrive tool that's rate-limited so fixing a bug doesn't create an incident."

3.5 Fault Line 5: Shared Multi-Tenant Cluster vs Dedicated Queues#

ModelUtilizationIsolationOperational CostWho Pays
One shared cluster, no quotasHighest (~60–70%)NoneOne team, one clusterThe quiet tenant, when the noisy one fills disks
Shared cluster with quotasHigh (~50–60%)Rate and storage, not I/O or failureOne team + quota managementTenants hitting quotas during legitimate spikes
Cells (N shared clusters, tenants assigned)Medium (~40–50%)Failure isolation per cell; blast radius = 1/NFleet tooling, placement servicePlatform team building cell routing
Dedicated cluster per teamLow (~15–30%)FullN × on-call, N × upgradesThe org: ~3–5× infra cost, and N teams who don't want to run brokers

The Staff default: a shared cluster per cell with per-tenant quotas on produce bytes/s, consume bytes/s and stored bytes; critical-tier tenants placed in a dedicated cell; best-effort tenants share. Start with one cell; split when a cell exceeds ~30% of org traffic or any tenant exceeds ~20% of a cell.

When to deviate: a single tenant that is 50%+ of total traffic (often the analytics firehose) should get its own cluster — not for isolation from others, but so others are isolated from it.

🎯 Staff Move: "Quotas protect capacity; cells protect against failure. A tenant that sends a 900 KB message with a pathological key can't be stopped by a rate limit — only by putting them somewhere their blast radius is bounded."


4. Failure Modes & Operational Reality#

4.1 The Lease-Expiry Redelivery Storm#

A downstream dependency slows down; handlers that took 5 s now take 70 s; the 60 s lease expires on every in-flight message; each is redelivered to another worker; the dependency's load doubles; handlers slow further.

t=0:        Downstream payments API p99 rises from 4s to 70s (its DB failover).
t=+60s:     Leases on ~12,000 in-flight messages expire. All redelivered.
            Workers now run 2 copies of each slow call. Downstream QPS ×2.
t=+2min:    Downstream p99 → 120s. Second wave of expiries. In-flight ×3.
            queue.redelivery_rate: 0.1% → 64%. ack.stale_receipt_rate spikes.
t=+3min:    Page: queue.redelivery_rate > 10% for 2m. Also downstream 5xx > 20%.
t=+5min:    On-call pauses consumption on the queue (dispatcher flag), not the producers.
t=+12min:   Downstream recovers. Consumers resumed at 25% concurrency, ramp 25%/2min.
t=+30min:   Backlog of 480K drained. 3,100 duplicate side effects — all deduplicated
            by consumer idempotency table except 41 from a handler without it.

Detection: queue.redelivery_rate, ack.stale_receipt_rate, handler.duration_p99 vs lease_seconds, downstream error rate.

Mitigation: heartbeats so live handlers keep their leases; a per-queue concurrency cap that the consumer SDK enforces; client-side circuit breaker on the dependency so handlers fail fast and nack with backoff instead of hanging.

Prevention: alert when handler.duration_p99 > 0.5 × lease_seconds; make heartbeat the SDK default; idempotency conformance test for every handler on a critical-tier queue.

Owner: consuming team (handler and dependency behavior); platform team owns the redelivery-rate alert and the pause control.

4.2 Poison Message Stalls an Ordered Partition#

t=0:        A producer deploy emits an event with a new enum value. Consumer v41
            throws on parse. Queue is ordered per account, partition-blocking.
t=+1s:      Partition 23 (≈ 9,000 accounts) stops advancing. Retries: 10s, 1m, 10m...
t=+5min:    queue.oldest_unacked_age on partition 23: 5 min. Other partitions fine.
            Aggregate lag looks normal — 1 of 64 partitions is stuck.
t=+40min:   Customers on partition 23 report stale balances. Support escalates.
t=+50min:   On-call finds the stuck offset via per-partition age dashboard.
t=+55min:   Manual skip of the message (no key-parking) → account 77812 state now
            inconsistent: later events applied without the skipped one.
t=+3h:      Consumer v42 with the new enum value deployed. Skipped message replayed by hand.

Detection: queue.oldest_unacked_age per partition (max, not average); consumer.parse_error_rate; partition.commit_stalled_seconds.

Mitigation: key-parking — the poisoned key's messages move to a hold queue; the partition continues. Redrive after fix preserves order for the key.

Prevention: schema registry with compatibility checks at produce time (new enum values are a breaking change for strict parsers); consumers treat unknown enum values as a handled case; per-partition age alert at 5 min for ordered queues.

Owner: producing team owns the schema change; consuming team owns the parser; platform owns key-parking and the per-partition alert. The post-mortem question is why a breaking schema change shipped without a compatibility check — that's the producer's gate.

4.3 Backlog Outlives Retention — The Silent Expiry#

t=0 (Fri 18:00): Consumer team's deploy breaks auth to its database. Handlers nack everything.
                 Retries exhaust → DLQ. DLQ alert routes to a Slack channel, not a pager.
t=+2h:           DLQ depth: 1.1M. Primary lag growing 40K/min (nacks are slow).
t=+60h (Mon):    Team notices. Fixes auth. Redrives DLQ.
t=+60h:          Primary queue retention: 4 days. OK. DLQ retention: 4 days (copied default).
t=+96h (Tue):    Oldest messages in the primary expire before consumers catch up:
                 consumers drain at 3× produce rate, but 60h of backlog needs ~30h.
                 ~210K messages expire unprocessed. No error, no log line — just gone.

Detection: queue.oldest_unacked_age approaching retention × 0.5; queue.expired_unprocessed_count (must exist — many systems don't emit it); drain ETA = backlog ÷ (consume rate − produce rate).

Mitigation: extend retention on the live queue (a config change — make sure it's online); add consumer capacity if the downstream can take it; prioritize oldest-first draining.

Prevention: retention ≥ 3× the longest plausible outage plus drain time; DLQ retention > primary; expired_unprocessed_count > 0 pages the platform on-call for critical-tier queues; DLQ alerts route to pagers for critical tier.

Owner: consuming team for the outage; platform team for retention defaults and the expiry metric. A message that expires unprocessed is a data-loss incident, not a consumer-lag ticket.

4.4 Acked Messages Lost After Broker Failure#

t=0:        Storage node B1 (leader for 140 partitions) has a degrading disk.
            Followers fall behind; in-sync set for 31 partitions shrinks to {B1}.
            Cluster configured min_in_sync=1 "for availability" during a migration.
t=+10min:   Producers still get acks for those 31 partitions — single copy.
t=+14min:   B1's disk fails hard. Controller elects lagging followers (unclean election
            enabled to avoid unavailability). ~38 seconds of writes on 31 partitions gone.
t=+14min:   Consumers see offsets jump backward then forward; ~1.9M acked messages missing.
t=+3h:      A downstream team's reconciliation flags missing orders. Incident declared.

Detection: partition.isr_size < min_in_sync_target (should page), cluster.unclean_elections_total, consumer-side gap detection on producer sequence numbers.

Mitigation: none after the fact, beyond replaying from producers' outboxes where they exist. That's the point: the only fix is before.

Prevention: min_in_sync = 2 is non-negotiable on the critical tier and requires a change review to lower; unclean leader election disabled; under-replication alert at 2 min; producers on critical flows use the outbox pattern so the source of truth can replay.

Owner: platform team. The deeper owner is whoever approved min_in_sync=1 "temporarily" — durability config changes need the same review as schema migrations.

4.5 Consumer Group Rebalance Storm#

t=0:        Consumer fleet of 200 pods autoscales to 320 under load.
t=+0–90s:   Each join triggers a group rebalance. Eager rebalancing revokes all
            partitions from all members on every change. 120 joins → near-continuous
            stop-the-world. Throughput: 180K msg/s → 15K msg/s.
t=+3min:    Lag climbs, autoscaler adds more pods → more rebalances.
t=+6min:    Page: consumer.throughput dropped 80%, group.rebalances_per_min > 20.
t=+8min:    On-call freezes the autoscaler. Group stabilizes in ~45s.

Detection: group.rebalances_per_min, group.partitions_unassigned_seconds, throughput vs consumer count (more pods, less throughput = rebalance problem).

Mitigation: incremental/cooperative rebalancing so only moved partitions pause; static membership so restarts don't trigger rebalances; autoscaling in steps with a cool-down ≥ 2× rebalance time.

Prevention: for work queues, the lease-based dispatcher sidesteps group rebalancing entirely — consumers are stateless pollers and adding one costs nothing. This is a strong argument for the lease layer over raw consumer groups when consumer counts are volatile.

Owner: consuming team (autoscaling policy); platform team (rebalance protocol defaults).

4.6 The Noisy Tenant Fills the Disks#

t=0:        Analytics tenant ships a bug: produces each event 40× (retry loop without
            backoff), 1.2 GB/s against a 300 MB/s quota... except storage quota was unset.
t=+20min:   Rate quota caps them at 300 MB/s, but at 7-day retention they're writing
            180 GB/10min into shared brokers. Disk usage 61% → 78%.
t=+50min:   Brokers hit 85% disk. Log cleaner falls behind. Page: broker.disk_used > 85%.
t=+55min:   Platform on-call cuts tenant's retention to 6h and storage quota to 2 TB.

Detection: broker.disk_used_pct, tenant.bytes_stored vs quota, tenant.produce_bytes_rate z-score vs 7-day baseline.

Mitigation: per-queue storage quota; when hit, that tenant's produces get 429 — never cluster-wide back-pressure. Emergency retention reduction is a pre-approved runbook step for best-effort tiers.

Prevention: every queue must declare a storage quota at creation; default derived from rate × retention × 1.5; anomaly alert on tenant produce rate.

Owner: platform team (quotas, isolation); tenant team (the bug). The tenant's quota hit should page them, not the platform.

4.7 Operational Reality Matrix#

FailureDetection SignalBlast RadiusMitigationOwner
Lease-expiry stormqueue.redelivery_rate > 10%, ack.stale_receipt_rateOne queue + its downstreamPause consumption, heartbeat, ramp back 25% stepsConsuming team
Poison in ordered queuequeue.oldest_unacked_age per partition > 5 minOne key (with parking) or one partitionKey-parking, DLQ, schema fixProducer (schema) + consumer (parser)
Backlog expiryqueue.expired_unprocessed_count > 0, age > 50% retentionAll messages older than retentionExtend retention live; drain oldest-firstPlatform (defaults) + consumer
Acked-then-lostpartition.isr_size < 2, unclean_elections_totalPartitions led by failed nodeNone after the fact; outbox replayPlatform team
Rebalance stormgroup.rebalances_per_min > 20One consumer groupFreeze autoscaler; cooperative rebalancingConsuming team
Noisy tenanttenant.bytes_stored vs quota; disk > 85%CellTenant quota 429s; emergency retention cutPlatform (isolation), tenant (bug)
Metadata store loss of quorummetadata.leader_present == 0Whole cell: no leader changes, no new queuesExisting leaders keep serving; restore quorumPlatform team
Dispatcher failoverdispatcher.lease_state_rebuild_secondsPartitions owned by that dispatcher: duplicates for in-flight leasesRebuild from ack log; in-flight redeliveredPlatform team
DLQ unowneddlq.oldest_message_age > 1h with no ack from ownerThat queue's failed messagesEscalate to owning team's manager after 24hConsuming team

🎯 Staff Insight: Notice that half the owners are consuming teams. A queue platform that pages itself for consumer problems burns out its on-call within a quarter. The platform's job is to make the signals precise and route them to the team that can fix them.


5. Evaluation Rubric#

5.1 Level-Based Signals#

DimensionSenior (L5)Staff (L6)Principal (L7)
FramingPicks a technology and describes itSeparates work queue / event log / fan-out and commits to oneAsks how many queueing systems exist and which should be retired
DurabilityRF=3Ack after cross-AZ quorum; rejects writes below min in-syncDefines durability tiers with $/GB and loss bounds; tier choice is a reviewed decision
Delivery"Exactly-once"At-least-once + idempotent effect via message_id in the consumer's transactionShips dedup in the SDK, enforces via readiness review and redelivery chaos tests
OrderingOne partition for orderPer key, one in-flight per key, key-parking for poisonTreats ordering as a priced product feature; pushes versioned events instead
FailureRetries + DLQ boxBounded retries, classified errors, owned DLQ with age alert and redriveTracks "messages expired unprocessed" org-wide as an SLO-class metric
Scale"Add partitions"Sizes partitions for 3× peak; knows repartitioning breaks key orderPlans cells, chargeback and a migration path off legacy brokers
OwnershipPlatform owns the queueConsumers own handlers and DLQs; platform owns contract, tooling and signalsWrites the standard; decides which use cases the platform refuses

5.2 Strong Hire Signals#

SignalWhat It Sounds Like
Defines the ack precisely"An ack means two replicas in two AZs have it. Below that, the producer gets an error, not a lie."
Treats duplicates as routine"Every consumer crash between side effect and ack is a duplicate. That's not an edge case; it's Tuesday."
Separates crash detection from job duration"60-second lease to notice a dead worker; heartbeats so a 10-minute job keeps its claim."
Bounds the blast radius of poison"A poison message should cost one key, not one partition."
Links retention to operations"Retention is our longest plausible consumer outage times three, plus drain time."
Routes failures to the right owner"The platform alerts; the consuming team gets paged for its own DLQ."

5.3 Lean No-Hire Signals#

SignalWhy It Misses the Bar
"Kafka gives us exactly-once" with no qualificationConfuses a transactional log feature with end-to-end effect on external systems
Global ordering by defaultCaps throughput and creates a single point of head-of-line blocking without asking whether order is needed
No answer for consumer crash mid-processingThe defining failure of a queue was never considered
DLQ mentioned without owner, alert or retentionThe design has a place where messages silently die
Spends 15 minutes on segment files and indexesRedesigning a solved storage layer instead of the delivery contract
Autoscale consumers as the answer to every backlogIgnores the downstream; more consumers often make the outage worse

5.4 Common False Positives#

  • Deep Kafka internals knowledge ≠ queue design. Reciting ISR mechanics, log.segment.bytes and the KRaft migration is impressive and orthogonal. The question is what the consumer is promised.
  • Drawing many retry topics ≠ owning failure. A 6-tier retry topology with no owner for the last tier is an elaborate way to delete data later.
  • Naming every broker ≠ judgment. A comparison table of RabbitMQ, Pulsar, SQS and Kafka is useful only after the contract is defined.
  • High throughput numbers ≠ correct sizing. "Kafka does millions per second" says nothing about whether one consumer can keep up with one partition.

6. Interview Flow & Pivots#

6.1 Typical 45-Minute Shape#

PhaseTimeGoal
Framing and intents0–4 minCommit to work queue; capture rate, size, processing time, outage tolerance
Entities and API4–7 minMessage ID, receipt, lease, ack/nack/extend
High-level architecture7–12 minFront-end, log, dispatcher, metadata — no more
Ack path12–20 minQuorum, min in-sync, producer idempotence
Consumer contract20–30 minLeases, heartbeats, idempotent effect, redelivery
Poison + DLQ + ordering30–37 minRetry tiers, key-parking, ownership
Multi-tenancy / scale37–42 minQuotas, cells, partition sizing
Wrap-up42–45 minContract summary, what's next, what you skipped

6.2 How Interviewers Pivot — And What They're Testing#

PivotWhat They're TestingStrong Response
"Now make it exactly-once."Whether you'll overclaimDefine exactly-once effect; show the consumer transaction; note where broker transactions do and don't help
"One customer needs strict global ordering."Pushback with costSingle-partition queue for that tenant, throughput ceiling in their SLO, or versioned events instead
"Consumers are 6 hours behind."Backlog reasoningAge vs retention, drain ETA, downstream capacity before adding consumers
"A whole AZ goes down."Durability mechanicsQuorum ack means no loss; leaders move in ~2–10 s; producers retry; consumers see brief redelivery
"Make it 10× bigger."Scale without re-architectureCells, partition headroom, replication bandwidth cost
"Messages are 50 MB videos."Payload boundariesClaim-check: blob storage + pointer; the queue carries metadata
"Delay a message by 3 days."Knowing what queues are bad atScheduler with a time index; the queue handles the "now" part

6.3 What to Deliberately Skip#

  • Segment file format, index files, zero-copy sendfile. One sentence and a link to the Kafka guide.
  • Leader election algorithm details. "Raft in the metadata group" suffices.
  • Wire protocol and serialization format. Mention schema registry for compatibility; skip Avro vs Protobuf debates.
  • Compression codec tuning. "Producer-side batch compression, ~3–5× on JSON" is enough.
  • Exact JVM/GC tuning of brokers. Operations detail with no design signal.

6.4 Follow-Up Questions to Expect#

  1. "What happens if the dispatcher holding lease state crashes?" — Rebuild from the compacted ack log; in-flight leases are treated as expired; duplicates possible, covered by idempotency.
  2. "How do you implement delayed delivery?" — Delay tiers (fixed-delay queues) for retries; a time-indexed store for arbitrary delays beyond ~15 min.
  3. "How do consumers dedupe without an unbounded table?" — Dedup window bounded by retention + max redelivery horizon; TTL the table at ~2× retention.
  4. "What if a producer sends to the wrong partition after you add partitions?" — For keyed queues, partition count changes require a migration: new queue, dual-write, drain, cut over.
  5. "How would you support priority?" — Separate queues per priority with weighted polling in the consumer SDK; not per-message priority inside a partition.
  6. "How do you bill tenants?" — Stored GB-hours + produced and consumed GB + request count; replication factor multiplies stored cost.
  7. "What would you change for a fan-out workload with 10,000 subscribers?" — Different system: drop per-message ack, use a log with many cheap readers or in-memory pub/sub.

7. Active Drills#

Drill 1: The Opening#

Prompt: "Design a distributed message queue."

Staff Answer

"There are three systems hiding in that sentence — a work queue where each job is done once and acked individually, an event log many teams replay, and ephemeral fan-out. I'll design the work queue as an internal multi-tenant platform: at-least-once, per-message ack, bounded retries into an owned DLQ, opt-in per-key ordering. I'll build on a partitioned replicated log so replay is available.

Assumptions: 200K msg/s peak, 2 KB average, handler p99 of 30 seconds with some 10-minute jobs, zero loss after ack, consumers tolerate duplicates, and we must survive a 72-hour consumer outage. Plan: define the producer ack → the consumer contract with leases → poison handling and DLQ ownership → tenancy and scale."

Why this is L6:

  • Names three intents with incompatible designs and commits to one
  • Asks for processing time and outage tolerance — the two numbers most candidates miss — because they drive leases and retention
  • Structures the interview around promises, not components

What L7 adds:

  • Asks which queueing systems already exist and whether this one replaces or joins them
  • Frames tenancy and ownership (who gets paged) as first-class requirements
  • Flags retention and partitioning of keyed queues as the long-lived decisions
❌ Common L5 Trap

"I'd use Kafka with a topic per use case, 3× replication, consumer groups for scaling, and exactly-once semantics enabled."

Why this misses: Every element is real, and none answers what happens when a consumer crashes after charging a card, or when one message can never be parsed. The interviewer's next three questions all land in territory the answer never touched.


Drill 2: The Ack#

Prompt: "Your producer got an ack. The leader's host dies 30 ms later. Is the message safe?"

Staff Answer

"It depends entirely on what the ack meant, so I'd define it up front: the ack is sent when the message is on the leader and at least one in-sync follower in a different AZ. With that policy, yes — the follower has it, gets promoted within ~2–10 seconds, and nothing acked is lost. If the policy were ack-on-leader, the 30 ms window is exactly the replication lag, and the message is gone.

Two supporting rules: if the in-sync set drops below two, produces are rejected — producers buffer and retry rather than receiving an ack backed by one copy; and unclean leader election is disabled, so an out-of-sync follower can never become leader and truncate acked data. The producer's retry of any in-flight batch is deduplicated by producer ID and sequence number."

Why this is L6:

  • Defines durability at ack time, in failure domains
  • Chooses visible unavailability over invisible loss, and says who absorbs it (producers' buffers)
  • Covers the retry duplicate that the failover itself creates

What L7 adds:

  • Offers tiered durability so analytics firehoses don't pay critical-tier latency and bandwidth
  • Puts min_in_sync changes under change review — the 4.4 incident starts with a "temporary" config change
  • Prices cross-AZ replication: at 400 MB/s and RF=3, ~800 MB/s cross-AZ, roughly 2 PB/month of billed transfer on clouds that charge for it

Drill 3: The Slow Consumer#

Prompt: "Most jobs take 200 ms, but some take 8 minutes. Consumers crash a few times a day. Pick a visibility timeout."

Staff Answer

"I won't pick one number, because it's two requirements: how fast we notice a dead worker, and how long a live worker may hold a job. Lease of 60 seconds for crash detection; the consumer SDK heartbeats every 20 seconds while the handler runs, extending the lease; a maximum total lease of one hour as a backstop against hung handlers with a live heartbeat thread. Every ack carries a receipt with a lease generation, so a worker whose lease expired can't ack a message now held by someone else.

Signals: handler.duration_p99 against the lease, ack.stale_receipt_rate, and redelivery rate. If stale receipts climb above ~0.1%, either heartbeats are failing or handlers are blocking the heartbeat thread."

Why this is L6:

  • Decouples crash detection from job duration with heartbeats
  • Bounds the failure mode of heartbeats themselves (hung handler, live heartbeat)
  • Makes stale acks impossible via receipt generations

What L7 adds:

  • Draws the boundary: anything over an hour should be a workflow with checkpoints, and the platform says so in its standard
  • Makes heartbeat a default in the SDK rather than an opt-in, because the teams who need it least will configure it most

Drill 4: Make It Concrete — Sizing#

Prompt: "200K msg/s peak at 2 KB. Size it."

Staff Answer

"400 MB/s ingress. Planning budget ~5 MB/s per partition so a single consumer instance can keep up and catch-up reads don't saturate a broker; at 3× headroom that's 1.2 GB/s ÷ 5 MB/s ≈ 240 partitions. Replication at RF=3 adds 800 MB/s of follower traffic; consumers read another ~400 MB/s, more during catch-up.

Storage: 400 MB/s × 86,400 s ≈ 35 TB/day raw, × 3 replicas ≈ 104 TB/day; 7-day critical retention ≈ 730 TB on disk before compression. With ~4× batch compression on JSON payloads, ~180 TB. At ~12 TB usable per broker that's ~16 brokers for storage alone; network usually binds first — at ~1 GB/s usable per broker for produce + replicate + consume, ~2 GB/s total traffic needs at least 6–8 brokers at 50% headroom. I'd start at 18 brokers across 3 AZs and offload segments older than 24 hours to object storage, which drops local disk by ~85% and lets broker count follow network instead of disk."

Why this is L6:

  • Derives partition count from per-consumer throughput, not broker capability
  • Counts replication and consumer traffic, not just ingress
  • Notices which resource binds first and uses tiered storage to fix the wrong one binding

What L7 adds:

  • Converts to dollars: cross-AZ transfer and storage dominate, not compute
  • Asks whether the 7-day retention is for the critical tier only — applying it to the best-effort firehose is the most expensive line item in the plan

Drill 5: The Dependency Goes Down#

Prompt: "Consumers write to a database that's now at 95% CPU. Lag is growing fast. The team wants to double the consumer count."

Staff Answer

"Doubling consumers doubles pressure on the thing that's already failing — it's the most common way a slowdown becomes an outage. The queue is doing its job: it's absorbing the backlog. The question is how long we can let it grow. Backlog age versus retention tells me: if the oldest message is 40 minutes old and retention is 7 days, we have days, not minutes.

Actions: cap consumer concurrency at the level the database sustained before the incident — the SDK reads a per-queue concurrency limit from config; have handlers fail fast with a circuit breaker and nack with backoff rather than hang; once the database recovers, ramp concurrency 25% every few minutes while watching its CPU. If the business needs some messages sooner, split by priority — that's a new queue, not a hack in this one."

Why this is L6:

  • Recognizes backlog as the queue working, not failing
  • Uses age vs retention to decide urgency
  • Applies backpressure at the consumer, with a ramp, instead of autoscaling into the outage

What L7 adds:

  • Proposes that consumer autoscaling be bounded by downstream capacity signals org-wide — a common failure across many teams
  • Ties it to an error budget: end-to-end latency SLO for this queue decides when lag becomes an incident

Drill 6: The Hot Key#

Prompt: "One tenant's queue is ordered by merchant_id. A single merchant is 30% of traffic. Their partition lags by hours."

Staff Answer

"With per-key ordering, that merchant is serial by definition — one in flight at a time — so no amount of partitions helps. First question: what invariant needs ordering per merchant? Often it's per-merchant-per-order, or per-account under the merchant. Narrowing the key to merchant_id:order_id keeps the invariant and spreads load across thousands of keys.

If ordering is genuinely per merchant, the options are: speed up the handler for that key (batch consecutive messages for the same key into one transaction — 50 messages per commit can be a 20–30× gain); or relax to versioned events where the handler drops stale versions and lets order be restored at read time. What I won't do is salt the key and pretend order is preserved."

Why this is L6:

  • Identifies that partitioning cannot fix a serial key
  • Narrows the ordering scope to the actual invariant
  • Offers batching per key as a throughput lever that preserves order

What L7 adds:

  • Points out the product conversation: the merchant SLA should reflect that their volume is serial work, or the contract should change
  • Uses this as evidence for a platform rule — ordering keys must be reviewed at queue creation, with a cardinality estimate

Drill 7: Multi-Tenant Isolation#

Prompt: "300 teams share the cluster. One team's bug just tripled their produce rate. How did your design protect everyone else?"

Staff Answer

"Three layers. Rate: per-tenant produce bytes/s and requests/s quotas at the stateless front end, enforced locally with token buckets synced every second — the tenant sees 429s, the cluster doesn't. Storage: a per-queue stored-bytes quota, so even at their allowed rate they can't fill shared disks beyond what rate × retention allows. Failure: critical-tier tenants live in a separate cell with their own brokers and metadata group, so the worst case — a pathological workload that a quota doesn't catch, like huge messages on a single key — stays within one cell.

The tenant's quota violation pages them. The platform gets paged only if a cell-level signal moves — disk > 85%, produce latency p99 across tenants > 2× baseline."

Why this is L6:

  • Separates rate, storage and failure isolation as different mechanisms
  • Routes the alert to the team that caused the problem
  • Uses cells for what quotas can't catch

What L7 adds:

  • Chargeback by stored GB-hours and transferred GB, so the tenant's bug shows up in their budget
  • Sets the cell-split policy (e.g., no tenant > 20% of a cell) as a standing rule rather than an incident follow-up

Drill 8: Build vs Buy#

Prompt: "Should we build this, run open-source Kafka, or use a managed cloud queue?"

Staff Answer

"Building a queue from scratch is almost never right — storage, replication and election are solved problems, and owning their bugs is a multi-year commitment. The real choice is managed queue vs self-run log vs managed log.

For work-queue semantics under ~10K msg/s per queue on one cloud, a managed queue like SQS is the default: per-message ack, leases, DLQ and redrive come built in, and nobody gets paged for brokers. For high-throughput event-log use with replay and many consumer groups, a log — managed if the team is under ~5 engineers dedicated to it. Self-running Kafka makes sense when the cost at scale (often 30–60% cheaper than managed at hundreds of MB/s) pays for a dedicated team of 4–6. What I'd build is the thin layer on top: the consumer SDK with idempotency and heartbeats, the DLQ tooling, and quotas — the parts that encode our contract."

Why this is L6:

  • Rejects building the storage engine, with a reason
  • Gives thresholds for each choice
  • Identifies the layer that is worth building — the contract layer

What L7 adds:

  • Prices it: a self-run team of 5 is ~$1.25M/yr fully loaded; break-even against managed fees occurs somewhere in the hundreds of MB/s
  • Considers lock-in: a managed queue's API is easy to wrap; a managed log's ecosystem (connectors, stream processors) is harder to leave

Drill 9: Changing a Keyed Queue Without an Outage#

Prompt: "A keyed, ordered queue has 32 partitions and needs 128. How do you do it live?"

Staff Answer

"Adding partitions in place changes hash(key) mod N for most keys, so a key's new messages land on a different partition while its old messages are still unprocessed elsewhere — ordering breaks during the transition. I'd do it as a migration:

  1. Create the new queue with 128 partitions.
  2. Producers dual-write — or better, flip producers via config to the new queue at a known point.
  3. For ordered keys, the dispatcher on the new queue holds a key's messages until the old queue has drained that key — or, simpler, until the old queue is fully drained, if drain time is minutes.
  4. Consumers read old-then-new; once the old queue's lag is zero, retire it.

If the dispatcher maps keys to partitions through a lookup table rather than modulo, I can move keys in small batches with per-key fencing instead — that's worth building if partition changes are frequent."

Why this is L6:

  • Knows that repartitioning silently breaks per-key ordering
  • Uses drain-then-cutover per key to preserve order
  • Identifies the indirection (key→partition map) that makes this a two-way door

What L7 adds:

  • Makes partition count a sized-for-two-years decision at queue creation, with platform-provided guidance
  • Turns the migration into self-service tooling so teams don't do it by hand at 2 a.m.

Drill 10: Multi-Region#

Prompt: "We're adding an EU region. Some queues must stay in the EU; others need to survive a full region loss."

Staff Answer

"Two different requirements. Residency: EU tenants' queues live only in EU cells; the front end routes by tenant home region, and cross-region produce from a US service to an EU queue is an explicit, audited path. Region-loss survival: synchronous cross-region replication would add 70–150 ms to every ack, so I'd use asynchronous mirroring to a standby region with an RPO of seconds and documented duplicate and loss bounds.

Failover is the hard part: consumer progress must be translated, because offsets differ between mirrored logs. I'd mirror ack state alongside messages and accept that, after failover, consumers see a window of redeliveries — covered by idempotency. Ordering per key holds within a region; during failover, the old region is fenced before the new one accepts writes for the same keys."

Why this is L6:

  • Separates residency from disaster recovery
  • Rejects synchronous cross-region acks with a latency number
  • Names offset translation and fencing as the real failover problems

What L7 adds:

  • Classifies which queues actually need regional DR — usually under 10% — and refuses to mirror the rest, saving the replication bill
  • Runs a twice-yearly regional failover game day with measured RPO and duplicate count

8. Deep Dive Scenarios#

Deep Dive 1: Peak-Traffic Incident — Black Friday Fulfillment Backlog#

Context: On Black Friday, produce rate on the fulfillment.jobs queue hits 9× normal. Lag reaches 2.4M messages; the oldest is 38 minutes old. The fulfillment team has already scaled consumers from 80 to 400 pods, and lag is growing faster. The warehouse-management database behind the consumers is at 98% CPU. The on-call escalates to you.

Questions to Surface First:

  • What's the consumption rate now vs at 80 pods? (If it went down, consumers are the problem.)
  • What's the redelivery rate? Are leases expiring because handlers slowed?
  • What's the oldest message's age vs retention and vs the business deadline (orders must ship same day if placed before 14:00)?
  • Is all work equal, or are some jobs (express shipping) more urgent?

Typical L5 Approach: Keeps adding consumers or asks the DBA to scale the database vertically; considers raising the visibility timeout to cut duplicates.

Staff Approach: Recognizes the death spiral: 400 pods overload the database, handlers slow past the 60 s lease, redelivery doubles load. Cuts concurrency back to the level the database sustained (~120 pods), confirms heartbeats are on, and splits express-shipping jobs into a priority queue so the business deadline is met for the work that has one.

Principal Approach: Treats "autoscale consumers on lag" as an org-wide anti-pattern. Proposes that consumer autoscaling policies be bounded by a downstream-capacity signal and that peak-season readiness reviews include a load test of the queue and its consumers' dependencies at 10×.

Staff Approach — Full Reasoning
PhaseWhat to Do
Immediate (0–5 min)Freeze consumer autoscaler. Set per-queue concurrency cap to 120 via config. Check queue.redelivery_rate — it's 41%, confirming lease-expiry amplification.
TriageDatabase CPU drops from 98% to 70% within 3 min. Consumption rate rises from 2.1K/s to 5.8K/s — higher than at 400 pods. Backlog ETA at current produce rate: ~3.5 h.
Quick fixProducer-side change behind a flag: express-shipping jobs go to fulfillment.jobs.express with a dedicated 40-pod pool. Express backlog drains in 12 min.
GuardrailsAlert handler.duration_p99 > 0.5 × lease_seconds. Autoscaler max bound by DB connection budget (max_pods = db_conn_limit / conns_per_pod).
Post-mortemWhy did the autoscaler's only input ignore downstream health? Why was there no priority separation for work with a deadline?

Metrics to Watch: queue.oldest_unacked_age, queue.consume_rate, queue.redelivery_rate, db.cpu_utilization, handler.duration_p99

Organizational Follow-up: fulfillment team owns concurrency limits for its queue; the platform publishes a "downstream-aware autoscaling" recipe; peak readiness checklists add consumer-dependency load tests.

Ownership Question: "Who decides the concurrency limit on a queue's consumers?" Staff answer: The consuming team, because they own the downstream. The platform provides the knob and a default, and alerts when lease expiries indicate the limit is wrong.

Key Takeaway: "Backlog is the queue doing its job. Lag is only an incident when its age threatens retention or a business deadline — and adding consumers to a saturated dependency turns lag into an outage."

What clears the Staff bar:

  • Identifies lease-expiry amplification as the mechanism making it worse
  • Reduces concurrency instead of increasing it, and proves it with consumption rate
  • Separates deadline-bound work into its own queue

Deep Dive 2: Silent Failure — The DLQ Nobody Owned#

Context: A product manager asks why ~4% of new users over the past two weeks never received a welcome email. No alert fired. You find the email.welcome DLQ holds 380,000 messages, oldest 13 days old. DLQ retention is 14 days.

Questions to Surface First:

  • Who owns this DLQ? Where does its alert route — pager, ticket, or nowhere?
  • What's the failure reason distribution? One bug or many?
  • When does the oldest message expire — tomorrow?
  • Are the messages still valid to process (is a 13-day-old welcome email harmful)?

Typical L5 Approach: Redrives the entire DLQ immediately; 380K emails go out in 10 minutes, the email provider throttles the account, and some users get welcome emails two weeks late.

Staff Approach: First extends DLQ retention to stop expiry tomorrow. Samples failures: 99% are a template-render error introduced by a deploy 13 days ago. Gets product sign-off on what to do with stale messages (welcome emails older than 3 days → send a different "getting started" email or drop). Fixes the template, redrives at a rate the email provider tolerates.

Principal Approach: Asks how many other DLQs are in the same state. Commissions a one-time audit of every DLQ in the org, and makes "DLQ has an owner and a paging alert" a precondition for creating any queue in the critical or standard tier.

Staff Approach — Full Reasoning
PhaseWhat to Do
Immediate (0–5 min)Extend DLQ retention from 14 to 21 days. Confirm no new messages expire.
TriageGroup by error: 376K TemplateRenderError: missing locale field; 4K transient SMTP 4xx. Deploy 13 days ago added a locale-dependent field.
Quick fixFix template. Product decision: messages > 3 days old get the generic onboarding email, not the welcome email. Redrive at 200/s with a filter on age.
GuardrailsAlert dlq.oldest_message_age > 1h routes to the email team's pager. dlq.depth added to the team's service dashboard.
Post-mortemWhy did the DLQ alert route to a channel nobody watched? Why did a template deploy not have a canary that checks error rates per locale?

Metrics to Watch: dlq.depth, dlq.oldest_message_age, dlq.inflow_rate vs queue.produce_rate (ratio > 0.1% is suspicious), redrive.success_rate

Organizational Follow-up: the email team owns the DLQ and its pager; the platform blocks creating standard/critical queues without a DLQ owner and alert route; quarterly report of DLQ age by team.

Ownership Question: "Who owns a dead-lettered message?" Staff answer: The consuming team. The platform guarantees the message is kept, visible and redrivable; only the consuming team knows whether it should be fixed, replayed, transformed or dropped.

Key Takeaway: "A DLQ without a paging owner is a 14-day delay on deleting data. The decision about a stale message is a product decision, not an infrastructure one."

What clears the Staff bar:

  • Stops the bleeding (expiry) before investigating
  • Involves product on what to do with stale messages instead of blindly redriving
  • Redrives at a rate the downstream can absorb

Deep Dive 3: Large-Customer Onboarding — The Ledger Team Wants Strict Order#

Context: The ledger team wants to move from a database-polling job system onto the platform. Requirements: 60K msg/s, strict order per account, zero loss, and "exactly-once." Their largest account (an internal treasury account) is 8% of traffic. They want go-live in six weeks.

Questions to Surface First:

  • What does "exactly-once" mean for them — no double posting? Their ledger is the effect; can postings be idempotent by posting_id?
  • Is order needed per account, or per account-and-currency? What breaks if two postings for one account arrive reversed?
  • How fast is a single posting for the treasury account? 8% of 60K = 4.8K/s on one key.
  • What's their tolerated end-to-end latency, and what's their outage tolerance?

Typical L5 Approach: Gives them a single-partition queue for strict order and enables broker transactions for exactly-once.

Staff Approach: Converts "exactly-once" into a unique posting_id constraint in the ledger's transaction. Discovers the treasury account needs ~4.8K postings/s serially — a ~0.2 ms budget per posting, impossible with one transaction each. Proposes per-key batching in the consumer: consecutive messages for the same account committed in one ledger transaction (up to 200), preserving order and dropping per-posting overhead. Places them in the critical-tier cell.

Principal Approach: Notices this is the platform's first critical financial tenant and the beginning of a tier that needs its own durability standard, audit trail on redrive and skip operations, and change controls. Uses the onboarding to define that tier, rather than special-casing one team.

Staff Approach — Full Reasoning
PhaseWhat to Do
Immediate (week 1)Replace "exactly-once" with an explicit contract: at-least-once delivery, idempotent postings by posting_id unique index. Agree on per-account ordering with key-parking.
Triage (design)Load test: one account, sequential, 1 txn per posting → ~900/s. With per-key batching of 200: ~6.5K/s. Meets 4.8K/s with headroom.
Quick fix (rollout)Shadow mode: dual-consume from old system and new queue, compare postings for 2 weeks, 0 mismatches required.
GuardrailsSkip and redrive on this tier require two-person approval and are logged to an audit topic. queue.oldest_unacked_age per key-parked account pages at 5 min.
Post-mortem (retrospective)Document the critical-financial tier: RF=3, min_in_sync=2, fsync every 100 ms, 14-day retention, audit on operator actions.

Metrics to Watch: queue.oldest_unacked_age (per partition), consumer.batch_size_p50, ledger.duplicate_posting_rejects, parked_keys.count

Organizational Follow-up: the platform publishes the critical-financial tier; the ledger team owns their handlers and parked keys; finance signs off on skip procedures.

Ownership Question: "If a posting is parked for an hour, who decides whether to skip it?" Staff answer: The ledger team's on-call, with a second approver, and only after finance confirms the downstream effect. The platform enforces the two-person rule and keeps the audit trail.

Key Takeaway: "Exactly-once is a property of the consumer's transaction. Strict order on a hot key is a throughput ceiling you must measure on day one — batching per key is how you raise it without breaking order."

What clears the Staff bar:

  • Translates "exactly-once" into a concrete idempotency mechanism
  • Computes the serial budget for the hot key and finds it infeasible before launch
  • Uses shadow consumption and an explicit cutover criterion

Deep Dive 4: Post-Mortem — 1.9M Acked Messages Lost#

Context: You're asked to lead the post-mortem for incident 4.4: a storage node failure lost ~38 seconds of acked writes on 31 partitions. Root cause, per the initial write-up: "disk failure." Leadership wants to know whether it can happen again.

Questions to Surface First:

  • Why was min_in_sync set to 1? Who changed it, when, and was it reviewed?
  • Why was unclean leader election enabled?
  • How long were those partitions under-replicated before the failure? Did an alert fire?
  • Which consumers or producers could detect the gap? Can any data be recovered from producer outboxes?

Typical L5 Approach: Root cause: disk failure. Action items: replace the disk vendor, add disk SMART monitoring.

Staff Approach: Disk failure is the trigger, not the root cause — disks fail weekly in a fleet this size. Root cause: a durability config was lowered during a migration three weeks earlier and never restored, unclean election was enabled, and the under-replication alert was a dashboard, not a page. Actions: restore and lock config, page on under-replication, recover what's possible from producers' outboxes.

Principal Approach: Classifies durability config as a safety control. Puts it under policy-as-code with drift detection, requires a time-boxed exception with an owner and expiry for any reduction, and adds a monthly automated audit that compares live config to the tier standard.

Staff Approach — Full Reasoning
PhaseWhat to Do
Immediate (0–5 min)Audit all partitions for min_in_sync < 2 and unclean.election = true. Fix any found.
TriageChange log: min_in_sync=1 set 21 days ago for a broker migration "for availability"; no expiry. Under-replication on B1's partitions began 10 min before failure; alert was dashboard-only.
Quick fixProducers with outboxes (orders, payments) replay the window — 1.4M recovered. Remaining 0.5M (analytics, notifications) declared lost with product sign-off.
GuardrailsConfig drift detector: live config vs tier standard every 5 min, pages on mismatch. partition.isr_size < 2 pages after 2 min. Unclean election hard-disabled cluster-wide.
Post-mortemFive whys ends at: durability config was treated as an operational knob, not a safety control.

Metrics to Watch: partition.isr_size, partition.under_min_isr_seconds, config.drift_count, cluster.unclean_elections_total

Organizational Follow-up: platform owns durability config with change review; migrations that require lowering it need an exception with expiry and a named owner; producers on critical flows must use an outbox (recoverability as a requirement, not luck).

Ownership Question: "Who is accountable for the 0.5M messages that couldn't be recovered?" Staff answer: The platform team — we acked them. We notify each affected tenant with exact partition and time ranges so they can reconcile; that notification is part of the incident, not an afterthought.

Key Takeaway: "The disk was the trigger. The root cause was a durability setting that could be lowered without review and never came back. Safety config needs the same rigor as schema migrations."

What clears the Staff bar:

  • Refuses "hardware failure" as a root cause
  • Finds the config drift and the missing page
  • Recovers data via producer outboxes and owns notifying tenants

Deep Dive 5: Multi-Region Expansion — EU Launch With Residency#

Context: The company is launching in the EU in two quarters. Legal requires EU customer data to stay in the EU. Two critical-tier tenants (payments and orders) also want to survive a full region outage. Today everything runs in one US region.

Questions to Surface First:

  • Which queues carry EU personal data? Does the residency rule cover message payloads only, or metadata (keys, headers) too?
  • For DR: what RPO and RTO do payments and orders need? Is seconds of loss acceptable or must it be zero?
  • Do any flows legitimately cross regions (EU orders triggering US fulfillment)?
  • Who decides when to fail over a region — automation or a human?

Typical L5 Approach: Deploys a second cluster in the EU and mirrors everything bidirectionally.

Staff Approach: Splits residency from DR. EU tenants get queues in an EU cell; the front end routes by tenant home region. Cross-region flows go through an explicit, audited bridge that strips or tokenizes personal data. For DR, payments and orders mirror asynchronously to a standby region within the same jurisdiction, with ack state mirrored for offset translation; RPO measured at ~2–5 s.

Principal Approach: Uses the launch to introduce a region-and-tier placement model: every queue declares data classification and DR tier at creation, and the platform enforces placement. Calculates the mirroring bill and limits DR mirroring to queues whose owners accept the cost on their budget.

Staff Approach — Full Reasoning
PhaseWhat to Do
Immediate (planning)Inventory queues by data classification. Of 5,000 queues, ~600 carry EU-scoped personal data; 40 need DR.
Triage (design)EU cell with its own metadata group. Region bridge service for approved cross-region flows, with payload filtering and an allowlist.
Quick fix (DR)Async mirror for the 40 DR queues to a second EU region. Mirror ack state; fence the primary before promoting the standby.
Guardrailsmirror.lag_seconds alert at 10 s. Quarterly region failover game day; measure RPO and duplicate rate.
Post-mortem (readiness)Document: synchronous cross-region ack rejected (+70–150 ms per produce); RPO of seconds accepted by payments with outbox replay as backstop.

Metrics to Watch: mirror.lag_seconds, bridge.messages_blocked_by_policy, failover.duplicates_delivered, region.produce_latency_p99

Organizational Follow-up: legal and security own the data-classification policy; the platform enforces placement; tenant teams declare classification at queue creation.

Ownership Question: "Who presses the button to fail over a region?" Staff answer: The incident commander, following a runbook with explicit criteria — not automation — because failing over with seconds of RPO is a business decision about acceptable loss.

Key Takeaway: "Residency and disaster recovery are two different requirements with two different mechanisms. Mirror only what has an owner willing to pay for it."

What clears the Staff bar:

  • Separates residency from DR and solves each explicitly
  • Rejects synchronous cross-region acks with a latency cost
  • Handles offset translation and fencing as first-class failover problems

9. Level Expectations Summary#

After studying this case study, you should be able to:

  • Distinguish work queue, event log and ephemeral fan-out, and explain why one design can't serve all three well
  • Define the producer ack in terms of replicas and failure domains, and defend rejecting writes below the in-sync minimum
  • Explain why delivery is at-least-once and how consumers achieve exactly-once effect
  • Design leases with heartbeats and receipt generations, and size the lease from crash-detection needs
  • Bound the blast radius of a poison message with per-key in-flight limits and key-parking
  • Set retention from the longest plausible consumer outage plus drain time
  • Assign ownership: platform owns contract, tooling and signals; consuming teams own handlers and DLQs
  • Size partitions and brokers from throughput, replication and retention — and say which resource binds first
  • Explain why repartitioning a keyed queue is a migration, not a config change

The Bar for This Question#

Mid-level (L4): Describes producers, brokers and consumers correctly; picks a well-known product; mentions replication and consumer groups. Can explain what a queue is for. Struggles when asked what happens on consumer crash.

Senior (L5): Designs a reasonable partitioned, replicated system with consumer groups, retries and a DLQ. Knows the terms — ISR, offsets, visibility timeout. Handles the happy path and the obvious failures well. Under pressure, claims exactly-once, defaults to global or partition ordering, and treats the DLQ as a box on the diagram.

Staff+ (L6): Starts from the contract. Defines the ack in failure domains, the consumer's obligations, the lease and its heartbeat, the poison path and its owner. Treats duplicates and backlogs as routine operation, not edge cases. Assigns failures to the teams that can fix them. Uses numbers to make decisions — partition budget, lease length, retention — and knows which decisions are hard to reverse. The interviewer should learn something from the answer.


10. Staff Insiders: Controversial Opinions#

10.1 "Exactly-Once Delivery" Is the Wrong Requirement#

ClaimEvidence
No queue can deliver exactly once across an external side effectThe consumer can crash after the effect and before the ack; the broker cannot know which happened
Broker transactions solve a narrower problemThey make read-from-log, write-to-log atomic; an email, a card charge or a database write in another system is outside the transaction
Idempotent effect is cheapA unique index on message_id in the consumer's own database turns duplicates into no-ops at ~one extra row per message

The Staff position: Ask for at-least-once delivery and exactly-once effect. Put dedup where the side effect commits.

Why this matters in interviews: "Exactly-once" is the single most common overclaim in queue interviews. Correcting it — politely and with the mechanism — is a strong Staff-level signal.

10.2 Most Teams Asking for Ordering Don't Need It#

ClaimEvidence
Ordering is usually a proxy for "don't apply stale state"Versioned events with a "drop if older" rule achieve the same correctness without blocking
Ordering has a real costThroughput per key is serial; poison messages block; repartitioning becomes a migration
Ordering is rarely end-to-end anywayProducer retries, multiple producers and clock skew reorder events before they reach the queue

The Staff position: Default to unordered. Grant per-key order when a named invariant requires it, and make the requester own the throughput ceiling.

Why this matters in interviews: Pushing back on a requirement with a cheaper mechanism that preserves the invariant is exactly what Staff engineers do in design reviews.

10.3 Your Queue's Biggest Risk Is Its Consumers, Not Its Brokers#

ClaimEvidence
Broker failures are well handled by mature systemsLeader election and replication are decades-old solved problems in production logs
Most queue incidents start in consumersLease storms, poison messages, rebalance storms, unowned DLQs and autoscaling into dependencies all originate consumer-side
Platform teams invest in the wrong layerBroker tuning is visible and satisfying; consumer SDK defaults are where reliability is decided

The Staff position: Spend the platform's engineering budget on the consumer SDK — heartbeats, idempotency, concurrency limits, error classification — before the next broker optimization.

Why this matters in interviews: It shifts the conversation to where real incidents happen, which interviewers who've been on call recognize immediately.

10.4 A Database Table Is the Right Queue More Often Than You Think#

ClaimEvidence
Most internal job queues are smallMany services enqueue well under 1K jobs/s
Transactional enqueue removes a whole failure classThe job commits with the business row; no outbox, no dual write
SKIP LOCKED makes polling contention manageableWorkers grab distinct rows without blocking each other

The Staff position: Below ~1K jobs/s, when the job must be consistent with business data, a Postgres table is a fine queue. Graduate to a log when throughput, fan-out or retention needs outgrow it.

Why this matters in interviews: Knowing when not to introduce infrastructure is a judgment signal; interviewers notice candidates who reach for a cluster to move 50 jobs a second.

10.5 DLQ Retention Should Be Longer Than Primary Retention#

ClaimEvidence
Messages reach the DLQ exactly when humans are slowestFailures cluster around deploys, incidents and weekends
DLQ storage is cheapDLQ volume is typically < 0.1% of primary volume
Copying the primary's retention is the common defaultWhich is how DLQs expire messages before anyone looks

The Staff position: DLQ retention ≥ 2× primary retention, with an age alert at 1 h. Expiry from a DLQ should be an incident.

Why this matters in interviews: It's a small, concrete, defensible choice that shows you've watched a DLQ expire.


11. The Principal Lens (L7)#

Why L7 Sees This Problem Differently#

The Staff engineer designs a queue with a correct contract. The Principal engineer notices that the org runs five queueing systems — a self-run Kafka for events, SQS for some teams' jobs, a RabbitMQ cluster from an acquisition, a Redis-based job runner, and a Postgres jobs table in the monolith — each with its own delivery semantics, its own DLQ convention and its own on-call. Incidents don't come from any single system's design; they come from the seams: a team that assumes Kafka redelivers like SQS, a DLQ convention that only exists in one of them, five sets of runbooks. The L7 problem is asynchronous-messaging governance: which systems the org keeps, what contract all of them honor, and who is accountable for messages between teams.

The Org-Level Fault Line#

One messaging platform vs a curated set of systems vs team choice.

OptionWhat WorksWhat BreaksWho Pays
Every team choosesSpeed; teams pick what they know5+ systems, N on-call rotations, inconsistent semantics at team boundaries, no shared toolingOn-call engineers (toil), the org (duplicated infra), customers (seam incidents)
One platform for everythingOne contract, one SDK, one on-callForces ephemeral fan-out and analytics firehoses into a durability model built for jobs; platform becomes a bottleneckTeams with mismatched workloads (cost, latency); platform team (every request is theirs)
Curated set: one log, one work queue, one pub/sub; shared contractEach intent gets a fitting system; SDK and conventions unify the seamsRequires a standard and enforcement; three systems still to runPlatform team (three systems, one contract); teams migrating off legacy

🧭 Principal Move: "Three systems, one contract. A log for events, a work queue for jobs, pub/sub for ephemeral fan-out — and every one of them honors the same consumer SDK, DLQ ownership rule and durability tiers. The acquisition's RabbitMQ and the Redis job runner get an 18-month retirement plan."

Cost Model#

Assumptions: cloud-hosted, three AZs, RF=3, average message 2 KB, 7-day retention on critical tier and 2 days on best-effort, tiered storage for segments older than 24 h, fully loaded engineer ~$250K/year, cross-AZ transfer billed at typical cloud rates.

ScaleThroughputInfra ($/month)HeadcountOn-call LoadDominant Cost Driver
Startup~2K msg/s, 4 MB/s~$1–3K managed queue/log0.5 eng (shared infra team)Shared rotation, < 1 page/monthManaged service fees; engineering time
Growth~50K msg/s, 100 MB/s~$30–60K (self-run log, 9–12 brokers, or managed at ~1.5–2× that)4–6 eng platform teamDedicated rotation, 3–6 pages/month, most consumer-causedCross-AZ replication transfer (~200 MB/s) and storage
Enterprise~1M msg/s, 2 GB/s across cells~$400K–900K (multiple cells, tiered storage, mirroring)15–25 eng (storage, SDK, tooling, tenancy)Per-cell follow-the-sun; consumer pages routed to tenantsReplication and mirroring transfer; retention on best-effort data

The pricing insight: at enterprise scale, best-effort tenants with 7-day retention they don't need are often 30–50% of the storage bill. A retention-tier policy with chargeback typically pays for two platform engineers in its first quarter. Conversely, cross-region mirroring for every queue "just in case" can double the transfer bill; mirroring only the ~5–10% of queues with a real DR requirement is the single largest lever.

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionDoor TypeReversibility Cost
Partition count and key for an ordered queueOne-way (without key→partition indirection)Migration with drain-and-cutover per queue; per-tenant coordination
Message ID and envelope format (headers, schema reference)One-wayEvery producer and consumer changes; dedup tables keyed on the old ID
Delivery semantics promised to consumersOne-wayWeakening a guarantee silently breaks consumers that relied on it
Exposing broker-native APIs vs a platform APIOne-way-ishTeams bind to broker clients; swapping brokers becomes N migrations
Choice of broker behind the platform APITwo-way (if the API abstracts it)Months of migration, but invisible to tenants
Lease length, retry schedule, DLQ thresholdsTwo-wayConfig change per queue
Retention per tierTwo-way upward; one-way downwardShortening retention deletes data immediately
Tiered storage on/offTwo-wayBackground segment movement

The Standard I'd Write#

RFC-MSG-001: Asynchronous Messaging Standard
Status: Approved   Owners: Messaging Platform + Production Engineering

Scope
  Every service that produces or consumes messages across a team boundary,
  on any approved system (event log, work queue, pub/sub).

MUST
  1. Use an approved system: the event log, the work queue, or pub/sub.
     New use of other brokers requires an exception.
  2. Declare tier at creation: ephemeral, standard or critical. Tier sets
     ack policy, retention, DR eligibility and alert routing.
  3. Carry a stable message_id in the platform envelope. Consumers of standard
     and critical queues MUST be idempotent on message_id and pass the
     redelivery conformance test.
  4. Ack only after the side effect commits.
  5. Every standard and critical queue has a DLQ with a named owning team,
     a paging alert on oldest-message age (critical: 15 min, standard: 1 h),
     and DLQ retention at least 2x primary retention.
  6. Messages expired unprocessed on critical queues are incidents.
  7. Durability settings below the tier standard require a time-boxed
     exception with an owner and expiry.

SHOULD
  1. Prefer unordered queues; per-key ordering requires a stated invariant
     and a key cardinality estimate.
  2. Payloads over 256 KB use claim-check via blob storage.
  3. Use the platform consumer SDK (heartbeats, concurrency limits, error
     classification) rather than raw broker clients.

Exceptions
  Filed with Messaging Platform; reviewed within 5 business days; time-boxed
  to 2 quarters; durability exceptions need Production Engineering sign-off.

Success metrics
  - Acked-then-lost messages on critical tier: 0
  - Messages expired unprocessed (critical): 0; (standard): < 1 per million
  - DLQs with no owner: 0
  - Median DLQ resolution time: < 24 h
  - Legacy brokers in production: 0 by end of year 2

What I'd Tell the VP#

"We run five different messaging systems, and most of our async incidents happen in the gaps between them — a team assuming one system behaves like another, or failed messages piling up somewhere nobody watches. I'm proposing we consolidate to three systems that share one contract: what a 'sent' message guarantees, who owns failures, and how long we keep data. It's about four engineers for a year, retires two legacy clusters, and should cut async-related incidents by roughly half. We'll also introduce retention tiers and chargeback, which we expect to reduce the messaging infrastructure bill by 25–35%. The main risk is migration effort for teams on the legacy systems; we'll provide tooling and an 18-month window."

Principal Interview Signals#

SignalWhat It Sounds Like
Governs the class, not the instance"The question isn't this queue — it's why we run five, and which contract all of them honor."
Prices the tradeoff"Mirroring everything doubles our transfer bill; mirroring the 40 queues with real DR needs costs 6% of that."
Identifies one-way doors"The envelope format and message ID are the decisions I'd slow down on. Lease lengths we can change on Tuesday."
Routes ownership structurally"Consumer failures page consumer teams. A platform that pages itself for its tenants' bugs doesn't survive a year."
Knows when not to standardize"Ephemeral fan-out stays out of the durability standard — forcing it in would cost latency and money for nothing."

Staff answers that L7 interviewers find insufficient:

  • "I'd build a great work queue with leases and DLQs" — correct, but ignores the four other systems teams already use and the seams between them.
  • "Each team owns its DLQ" — names an owner but not the enforcement: creation-time checks, paging routes and an org-wide expiry metric.
  • "We'll replicate across regions for resilience" — no classification of which queues need it, no cost, no failover decision owner.

Appendices

Appendix A: Mechanics in Depth#

A.1 Lease-Based Delivery on Top of a Log#

The dispatcher turns an offset-based log into a per-message work queue. For each partition it owns, it keeps:

  • committed_offset — every message below this is acked
  • in_flight — map of offset → (lease_expiry, generation, consumer_id)
  • acked_above_commit — sparse set of acked offsets above the commit point
  • attempts — offset → attempt count (only for messages that have failed)
  • busy_keys — keys with a message in flight (for ordered queues)
on receive(consumer, max_n, lease_s):
    batch = []
    for offset in scan(from = committed_offset):
        if offset in in_flight or offset in acked_above_commit: continue
        msg = log.read(offset)
        if queue.ordered and (msg.key in busy_keys or msg.key in parked_keys): continue
        gen = next_generation(offset)
        in_flight[offset] = (now + lease_s, gen, consumer)
        if queue.ordered: busy_keys.add(msg.key)
        batch.append(msg with receipt = (partition, offset, gen))
        if len(batch) == max_n: break
    persist_lease_events(batch)              # to compacted ack log
    return batch

on ack(receipt):
    (p, offset, gen) = receipt
    if in_flight.get(offset).gen != gen: reject STALE_RECEIPT
    remove in_flight[offset]; acked_above_commit.add(offset); release key
    advance committed_offset while committed_offset in acked_above_commit
    persist_ack(offset)

on lease_expired(offset) or nack(offset, retry_after):
    attempts[offset] += 1
    if attempts[offset] > max_attempts:
        dlq.produce(msg, history); mark acked; if ordered: park(msg.key)
    else:
        delay_tier(attempts[offset]).produce(msg)  # 10s, 1m, 10m, 1h
        mark acked in this partition               # the retry copy carries the obligation

Why it's right: the log stays append-only and sequential; per-message state is only kept for the window between commit point and head (usually seconds of messages), so memory is bounded by in-flight count, not backlog size. Why it can go wrong: a single message stuck far below the head keeps acked_above_commit growing — cap the gap (e.g., 100K offsets) and force the stuck message to the retry path when exceeded.

A.2 Ordered Retry Without Head-of-Line Blocking#

For ordered queues, a retried message must not be overtaken by a later message for the same key. The dispatcher keeps the key in busy_keys until the retry resolves; later messages for that key are skipped during scans (they stay unacked in the log). If the retry exhausts attempts, the key is parked: its later messages are copied to a per-key hold area in order and acked in the main partition, so the partition's commit point advances. Redrive replays the hold area for that key, in order, before releasing the key.

A.3 Why Not Per-Message State Stored in a Database?#

A per-message row (status, lease_expiry) in a database is the simplest work queue, and fine at low rates. At 200K msg/s it means ~600K writes/s (insert, lease, delete) with index churn, and every receive is a range scan with locking. The log + lease layer keeps the per-message write path append-only and keeps mutable state proportional to in-flight messages.

Appendix B: Message Envelope and Identity#

envelope {
  message_id:     string   # producer-supplied UUIDv7 or derived from business key
  queue:          string
  key:            string?  # ordering and partitioning key, optional
  schema_ref:     string   # registry subject + version
  produced_at:    int64    # producer clock, ms
  enqueued_at:    int64    # broker clock, ms (authoritative for age)
  attempt:        int32    # set by dispatcher on delivery
  trace_context:  string   # propagated for end-to-end tracing
  headers:        map<string,string>
  payload:        bytes    # ≤ 256 KB; otherwise claim-check pointer
}

Identity rules. message_id must be stable across producer retries — generate it before the first send, not inside the retry loop. Derive it from a business key (order_id:event_type:version) when the producer itself may be re-run (batch jobs, replays), so a re-run produces the same IDs and dedup works across runs.

Consumer dedup table.

CREATE TABLE processed_messages (
  message_id   text PRIMARY KEY,
  processed_at timestamptz NOT NULL DEFAULT now()
);
-- in the handler's transaction:
INSERT INTO processed_messages (message_id) VALUES ($1)
  ON CONFLICT DO NOTHING RETURNING message_id;   -- no row returned → duplicate, skip effect
-- TTL: delete rows older than 2 × (retention + max retry horizon)

Appendix C: Coordination Mechanisms#

C.1 Producer Idempotence#

The broker assigns each producer session an ID; the producer numbers its batches per partition. The leader rejects a batch whose sequence number it has already appended (duplicate retry) and fails a batch with a gap (lost batch). This makes retries after a lost ack safe within a producer session. It does not cover a producer process restart that re-sends the same business event — that's what a business-derived message_id is for.

C.2 Quick Comparison#

MechanismProtects AgainstDoesn't Protect AgainstCost
Quorum ack (min_in_sync=2)Leader host/AZ loss after ackTwo-AZ correlated loss before flush~5–10 ms latency; 2× cross-AZ bandwidth
Producer idempotence (pid + seq)Duplicate append on retryProducer restart re-sendingNegligible
Broker dedup window (by message_id)Producer restarts within windowRe-sends after window (e.g., 5 min)Memory for window of IDs
Consumer dedup tableAll duplicate effectsHandlers with side effects outside the transactionOne row per message; TTL cleanup
Broker transactions (log → log)Duplicate output in read-process-write pipelinesExternal side effectsThroughput cost; transaction coordinator
Leases with generationsStale acks from expired consumersDuplicate effect from the expired consumer's completed workGeneration tracking
Outbox in producerLosing an event between DB commit and publish—Relay lag 100 ms–s

Appendix D: API Contract & Client Behavior#

BehaviorRule
Produce retriesExponential backoff from 50 ms, max 5 s, jitter ±50%; same message_id; producer buffers up to 60 s or 64 MB then fails the caller
429 Too Many RequestsHonor Retry-After; never retry immediately; surface quota exhaustion as a metric
ReceiveLong-poll up to 20 s; batch up to 10; never tight-loop on empty responses
AckAfter side effect commits; batch acks every 100 ms or 100 messages
NackWith retry_after for known transient errors; with permanent=true for validation errors → straight to DLQ
HeartbeatEvery lease/3; stop heartbeating if the handler exceeds max lease — let it be redelivered
ShutdownStop receiving, finish in-flight within grace period, nack the rest with zero delay
Thundering herd after outageConsumers resume with concurrency at 25%, ramp 25% per 2 min; producers flush buffers with jitter

Appendix E: Observability#

Core metrics (per queue, per partition where noted):

MetricWhy It MattersAlert
queue.oldest_unacked_age (per partition, max)The only lag metric that maps to user impact and retention riskCritical: > 5 min; standard: > 30 min
queue.expired_unprocessed_countSilent data lossAny non-zero on critical/standard
queue.redelivery_rateLease misconfiguration, consumer crashes, storms> 5% for 5 min
ack.stale_receipt_rateLeases shorter than real processing> 0.1%
dlq.oldest_message_ageUnowned or unnoticed failuresCritical: 15 min; standard: 1 h
partition.isr_sizeDurability below standard< 2 for 2 min pages
produce.latency_p99Ack-path health> 2× baseline for 5 min
tenant.quota_rejectionsTenant at limitRouted to tenant, not platform

Control plane vs data plane. The data plane (produce, receive, ack) must keep working when the control plane (metadata store, rebalancer, config service) is down: leaders keep serving, dispatchers keep leasing, front ends use cached partition maps. Control-plane outage blocks queue creation, config changes and leader moves — painful, but not an outage for tenants. Test this explicitly.

Debugging the silent failure. Aggregate lag hides single stuck partitions: always chart the max per-partition age, not the average. For "missing message" reports, trace by message_id through produce log → partition offset → lease events → ack or DLQ. If the platform can't answer "where is message X?" in under a minute, the incident will take hours.

Appendix F: Scale Evolution#

ScaleWhat WorksWhat Breaks Next
< 1K msg/sPostgres SKIP LOCKED table or a managed queueVacuum pressure, polling load
1K–50K msg/sManaged queue for jobs; managed log for eventsManaged costs; per-queue throughput caps on FIFO modes
50K–500K msg/sSelf-run log + dispatcher layer, one cell, quotasNoisy neighbors; single metadata group; cross-AZ bill
500K+ msg/sMultiple cells, tiered storage, placement service, chargebackCross-region needs, residency, migration tooling

Multi-region path: single region, three AZs → async mirroring for DR-tier queues to a standby region → region-homed tenants with residency enforcement → (rarely) active-active per key range with fencing. Most organizations never need the last step.

What You Don't Build on Day One: your own storage engine; cross-region active-active; per-message priority inside a partition; a custom delay scheduler for arbitrary future times (use a scheduler); broker-side transactions for external effects (they don't exist).

Appendix G: Multi-Tenancy, Fairness & Cost#

Fairness at the dispatcher. When many consumers poll a shared partition set, serve receives round-robin across queues with weights, not FIFO across all requests — otherwise a tenant with 500 polling workers starves a tenant with 5.

Quotas.

QuotaEnforcement PointOn Violation
Produce bytes/s, requests/sFront end, local token bucket, synced every 1 s429 with Retry-After
Consume bytes/sFront endThrottled receive (smaller batches)
Stored bytes per queueStorage tier, checked per segment rollProduce 429 for that queue only
Partitions per tenantControl plane at creationCreation rejected
In-flight leases per queueDispatcherReceive returns empty until acks arrive

Cost attribution. Charge on stored GB-hours × replication factor, produced and consumed GB, and mirrored GB. Show tenants their bill monthly; the first bill typically prompts 20–30% of tenants to shorten retention or move to the best-effort tier, which is the point.

  1. Loading the index…