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.
| Mode | Time | What to Read |
|---|---|---|
| Quick Review | 15 min | Executive Summary → Interview Walkthrough → Fault Lines table → Drills 1–3 |
| Targeted Study | 1–2 hrs | Executive Summary → Walkthrough → Section 3 (Fault Lines) → Section 4 (Failure Modes) → Deep Dives 1, 2 and 4 |
| Deep Dive | 3+ hrs | Everything, 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
| Model | How It Works | Pros | Cons |
|---|---|---|---|
| Per-message state queue (SQS, RabbitMQ classic) | Each message has its own state: visible, in-flight (leased), deleted. Consumers receive, then delete | Independent per-message ack; consumers scale freely; natural DLQ | Per-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 window | Very high throughput (sequential I/O); replay; many independent consumer groups | Parallelism 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 offsets | Replay and per-message ack; parallelism beyond partitions | More moving parts; ack-state must itself be durable |
Database-as-queue (SELECT … FOR UPDATE SKIP LOCKED) | Rows are jobs; workers lock and delete | Transactional with business data; zero new infra | Table bloat and vacuum pressure past ~1–5K jobs/s; polling load |
| In-memory broker (Redis lists/streams) | Messages in RAM, optional AOF | Sub-millisecond; simple | Durability 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#
| Behavior | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| First move | Draws producer → broker cluster → consumer; picks Kafka | Asks "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#
| Position | Rationale |
|---|---|
| Ack after quorum replication across failure domains, not after leader write | The 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 contract | Exactly-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 queue | Global 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 concurrency | The consumer knows its own capacity; push-based brokers have to guess and overrun slow consumers |
| Leases (visibility timeouts) with heartbeats, not fixed long timeouts | A 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 alert | Infinite retries convert one poison message into permanent load; an unowned DLQ converts it into silent loss |
| Retention ≥ 3× the longest plausible consumer outage | Backlog 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.
| Intent | Constraint | Strategy | Failure Mode | Correctness 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 isolated | Per-message leases with visibility timeout, bounded retries, DLQ, optional per-key ordering | Duplicate side effects on lease expiry; poison messages; backlog outliving retention | Zero 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 order | Partitioned replicated log, offsets per group, long retention or compaction | Hot partitions; consumer lag; schema breakage across teams | Per-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 worthless | In-memory pub/sub, no persistence or short TTL, drop on slow subscriber | Slow subscriber backs up the broker; message storms | Freshness 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 Line | The Tension |
|---|---|---|
| 1 | Ack Durability vs Producer Latency | Ack after leader write (~1–2ms) or after cross-AZ quorum (~5–15ms)? Who loses data when the leader dies in between? |
| 2 | Ordering vs Parallelism | Every ordering guarantee is a cap on concurrency and a head-of-line blocking risk. How much order is the business actually buying? |
| 3 | Redelivery Speed vs Duplicate Work | Short leases redeliver fast after a crash but duplicate slow work; long leases strand messages when consumers die |
| 4 | Retry Forever vs Dead-Letter | Keep retrying (no human needed, but poison loops) or give up after N (bounded load, but now someone must own the DLQ)? |
| 5 | Shared Multi-Tenant Cluster vs Dedicated Queues | Efficiency 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#
Quick-Reference: The 30-Second Cheat Sheet#
| Question | Staff 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#
| Number | Value | Context |
|---|---|---|
| Leader-only ack latency | ~1–2 ms | Page-cache write, same host; no durability across host loss |
| Cross-AZ quorum ack latency | ~5–15 ms p99 | One RTT to the fastest follower (~1–2ms) + batching + fsync policy |
| Per-partition throughput budget | ~5–10 MB/s planned, 30–50 MB/s ceiling | Plan at a fraction of the ceiling so a single consumer can keep up |
| Leader failover | ~2–10 s | Detection (session timeout) + election + client metadata refresh |
| SQS visibility timeout | 30 s default, 12 h max | The canonical lease parameter |
| SQS retention | 4 days default, 14 days max | Sets the outer bound on tolerable consumer outage |
| SQS long-poll wait | up to 20 s | Cuts empty receives (and their cost) by ~90%+ on idle queues |
| SQS FIFO dedup window | 5 min | Dedup state is bounded; producer retries happen within seconds |
Kafka min.insync.replicas for critical data | 2 (with RF=3) | Writes fail rather than ack with a single copy |
| Typical message size | 1–10 KB | Payloads > 256 KB–1 MB belong in blob storage with a pointer |
| Retry schedule | 10 s, 1 m, 10 m, 1 h, then DLQ | Bounded total ~1.2 h; covers most transient downstream outages |
| DLQ age alert | 1 h (critical: any message) | DLQ is a production queue, not an archive |
| Database-as-queue ceiling | ~1–5K jobs/s | Before vacuum/bloat and lock contention dominate on Postgres |
| Replication bandwidth | (RF − 1) × produce rate | 100 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:
| Question | Assumed Answer | Why It Matters |
|---|---|---|
| Peak produce rate? | 200K msg/s across all tenants, 2 KB average → ~400 MB/s | Sets partition count and replication bandwidth |
| Message size distribution? | p50 1 KB, p99 64 KB, hard cap 1 MB | Large payloads go to blob storage; changes batching |
| Processing time per message? | p50 200 ms, p99 30 s, some jobs 10 min | Decides lease length and heartbeat design |
| Ordering needed? | Some tenants need per-key order (account events); most don't | Ordering is opt-in per queue |
| Loss tolerance? | Zero after ack for critical tier | Ack after cross-AZ quorum |
| Duplicate tolerance? | Consumers must tolerate duplicates | At-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 queues | Quotas, 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)#
| Entity | Key Fields | Notes |
|---|---|---|
| Queue | queue_id, tenant, tier, partitions, ordering_mode, retention, max_attempts, dlq_id | Config is versioned; changes are rolled out like code |
| Message | message_id (producer-supplied or broker-assigned), key, payload, headers, enqueued_at, attempt | message_id is the idempotency handle for consumers |
| Lease | message_id, consumer_id, lease_expires_at, receipt handle | Receipt handle is per-delivery, so a stale consumer can't ack a re-leased message |
| Partition | partition_id, leader, replicas, high-watermark | Internal; 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.
"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).
"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).
"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#
| Mistake | Time Lost | Fix |
|---|---|---|
| Designing the log segment format and index files | 10+ min | Say "standard segmented log, see Kafka internals" and move on |
| Debating Kafka vs RabbitMQ vs SQS before defining requirements | 5 min | Define the contract first; the product choice follows |
| Drawing ZooKeeper election in detail | 5 min | "Metadata in a Raft group; leader election is a solved problem here" |
| Never reaching consumer failure or DLQ | Whole interview | Steer explicitly in Phase 4 — that's where the Staff-level signal is |
| Sizing partitions to three decimal places | 3 min | One 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#
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.
| Dimension | Work Queue | Event Log | Ephemeral Fan-out |
|---|---|---|---|
| Ack granularity | Per message | Per partition offset | None |
| Replay | Not required | Core feature | Never |
| Consumer scaling | Unbounded | ≤ partition count per group | Per subscriber |
| Slow consumer | Backlog grows; alert on age | Lag grows; alert on lag | Drop messages for that subscriber |
| Retention driver | Max consumer outage | Replay window, new consumers | Seconds or zero |
| Failure owner | Consuming team (DLQ) | Each consumer group | Nobody — loss is accepted |
2.2 When NOT to Use a Distributed Message Queue#
| Situation | Better Choice | Why |
|---|---|---|
| Caller needs the result to respond to the user | Synchronous RPC with timeout and retry | A queue adds latency and a second failure path, and the caller waits anyway |
| < 1K jobs/s and the job must commit with business data | SKIP LOCKED table in PostgreSQL | Transactional 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 Scheduling | Queues 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 compensation | Workflow engine (durable execution) — see Long-Running Processes | Chaining queues re-invents a state machine without visibility |
| Live updates where stale data is worthless | In-memory pub/sub or direct WebSocket fan-out | Durability adds latency and backlog with no value |
| "Decoupling" two services owned by the same team that deploy together | A function call | A 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 Assumption | Why It's Deliberately Vague | What to Say |
|---|---|---|
| Delivery semantics | To 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 scope | To see if you default to global order | "Ordering per key, opt-in. Most queues won't need it." |
| Processing time | To see if you notice leases depend on it | "What's the p99 handler time? Lease length follows from it." |
| Consumer outage duration | To see if you link retention to operations | "How long can a consumer team be down before we lose data? That sets retention." |
| Who owns failures | To see if you think organizationally | "The consuming team owns its DLQ; the platform owns the tooling and the alert." |
| Message size | To see if you push large payloads elsewhere | "Anything over 256 KB goes to blob storage; the message carries a pointer." |
| Number of tenants | To see if you consider isolation | "Shared cluster or per-team? That decides quotas versus cells." |
2.4 Precise Terminology#
| Term | Precise Meaning | Common 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 deleted | Acking on receipt, before processing |
| At-least-once | Every acked message is delivered ≥1 time; duplicates possible | Used as if duplicates were rare edge cases — they're routine on every crash |
| Exactly-once effect | Duplicates are delivered but have no additional effect, via idempotent consumers | "Exactly-once delivery," which no system provides across external side effects |
| Lease / visibility timeout | Time-bounded exclusive claim on a message | "Lock" — a lease expires without the holder's cooperation |
| Head-of-line blocking | A stuck message prevents later messages in the same ordering scope from being processed | Blamed on "slow consumers" when the cause is one poison message |
| Poison message | A message that fails deterministically regardless of retries | Confused with transient failure — retrying a poison message is pure waste |
| Consumer lag | Messages (or time) between the latest produced and the latest committed | Measured in messages only; age of oldest unacked message is what matters |
| Backpressure | Signal from a slower stage that limits the rate of a faster stage | "Autoscaling," which increases rate rather than limiting it |
| In-sync replica | A follower within the replication-lag bound that can be promoted without loss | Any 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 Policy | Producer p99 | What Survives | What Breaks | Who Pays |
|---|---|---|---|---|
| Fire-and-forget (no ack) | < 1 ms | Nothing guaranteed | Any network blip loses data silently | Downstream teams, who debug "missing events" for weeks |
| Leader write (page cache) | ~1–2 ms | Process crash on leader (page cache persists) | Leader host loss in the replication window: acked data gone; unclean election makes it worse | Customers whose messages vanish; the on-call who can't prove what was lost |
| Leader fsync | ~2–10 ms (SSD), batched | Leader host power loss | Leader disk/AZ loss — single copy | Same, less often; latency paid by every producer |
| Quorum: 2 of 3, cross-AZ | ~5–15 ms | Any single host or AZ loss | Simultaneous loss of two AZs before flush (very rare) | Producers pay ~5–10 ms; cluster pays (RF−1)× cross-AZ bandwidth |
| All replicas | Tail of the slowest replica, 20–100 ms+ | Same as quorum | One slow replica stalls every producer | Every 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 Scope | Max Parallelism | Poison Blast Radius | Who Pays |
|---|---|---|---|
| Global (one queue) | 1 consumer | Entire queue stops | Every consumer and every downstream user; throughput capped at ~10–50 MB/s forever |
| Per partition | = partition count | Every key on that partition (often 10K+ entities) | Unrelated customers who happen to hash to the same partition |
| Per key, partition-blocking | = partition count | Same as per partition — the common accidental design | Same; looks per-key on paper |
| Per key, key-parking | Unbounded (one in-flight per key) | One key | The one customer whose message is broken |
| None | Unbounded | One message | Consumers 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.
🎯 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 Strategy | Crash Recovery Time | Duplicate Rate | What Breaks | Who Pays |
|---|---|---|---|---|
| Short fixed (30 s) | ≤ 30 s | High for any job with p99 > 30 s | Slow jobs run 2–3× in parallel; downstream load multiplies | Downstream services (load), customers (duplicate emails/charges if not idempotent) |
| Long fixed (12 h) | ≤ 12 h | Near zero | Crashed consumer strands messages for hours; backlog invisible | Users waiting for work stranded in a dead worker |
| Short + heartbeat extension | ≤ lease (60 s) | Near zero for live consumers | Consumer that hangs while its heartbeat thread lives holds messages forever | Need a max total lease (e.g., 1 h) and handler-level watchdogs |
| Adaptive (lease = k × p99 handler time) | Varies | Low | Complex; p99 shifts during incidents exactly when you need stability | Platform 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.
| Policy | Load From Poison | Human Required | What Breaks | Who Pays |
|---|---|---|---|---|
| Immediate retry, unbounded | 100% of consumer capacity per poison message, forever | No | One bad message burns a worker permanently; storms during downstream outages | Consumer fleet; downstream dependency during its recovery |
| Backoff retry, unbounded | Small but permanent | No | Poison messages accumulate; retention eventually expires them silently | Customers whose messages expire; nobody sees it |
| Bounded retry → DLQ, unowned | Bounded | Nobody, in practice | DLQ fills; messages age out on day 14 | Same customers, with a false sense of safety |
| Bounded retry → DLQ, owned + alerted | Bounded | Yes, the consuming team | Toil if the failure rate is high | Consuming team's on-call — correctly, because they own the fix |
| Bounded retry → drop + log | Bounded | No | Data loss by design | Acceptable 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#
| Model | Utilization | Isolation | Operational Cost | Who Pays |
|---|---|---|---|---|
| One shared cluster, no quotas | Highest (~60–70%) | None | One team, one cluster | The quiet tenant, when the noisy one fills disks |
| Shared cluster with quotas | High (~50–60%) | Rate and storage, not I/O or failure | One team + quota management | Tenants hitting quotas during legitimate spikes |
| Cells (N shared clusters, tenants assigned) | Medium (~40–50%) | Failure isolation per cell; blast radius = 1/N | Fleet tooling, placement service | Platform team building cell routing |
| Dedicated cluster per team | Low (~15–30%) | Full | N × on-call, N × upgrades | The 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#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Lease-expiry storm | queue.redelivery_rate > 10%, ack.stale_receipt_rate | One queue + its downstream | Pause consumption, heartbeat, ramp back 25% steps | Consuming team |
| Poison in ordered queue | queue.oldest_unacked_age per partition > 5 min | One key (with parking) or one partition | Key-parking, DLQ, schema fix | Producer (schema) + consumer (parser) |
| Backlog expiry | queue.expired_unprocessed_count > 0, age > 50% retention | All messages older than retention | Extend retention live; drain oldest-first | Platform (defaults) + consumer |
| Acked-then-lost | partition.isr_size < 2, unclean_elections_total | Partitions led by failed node | None after the fact; outbox replay | Platform team |
| Rebalance storm | group.rebalances_per_min > 20 | One consumer group | Freeze autoscaler; cooperative rebalancing | Consuming team |
| Noisy tenant | tenant.bytes_stored vs quota; disk > 85% | Cell | Tenant quota 429s; emergency retention cut | Platform (isolation), tenant (bug) |
| Metadata store loss of quorum | metadata.leader_present == 0 | Whole cell: no leader changes, no new queues | Existing leaders keep serving; restore quorum | Platform team |
| Dispatcher failover | dispatcher.lease_state_rebuild_seconds | Partitions owned by that dispatcher: duplicates for in-flight leases | Rebuild from ack log; in-flight redelivered | Platform team |
| DLQ unowned | dlq.oldest_message_age > 1h with no ack from owner | That queue's failed messages | Escalate to owning team's manager after 24h | Consuming 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#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Framing | Picks a technology and describes it | Separates work queue / event log / fan-out and commits to one | Asks how many queueing systems exist and which should be retired |
| Durability | RF=3 | Ack after cross-AZ quorum; rejects writes below min in-sync | Defines 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 transaction | Ships dedup in the SDK, enforces via readiness review and redelivery chaos tests |
| Ordering | One partition for order | Per key, one in-flight per key, key-parking for poison | Treats ordering as a priced product feature; pushes versioned events instead |
| Failure | Retries + DLQ box | Bounded retries, classified errors, owned DLQ with age alert and redrive | Tracks "messages expired unprocessed" org-wide as an SLO-class metric |
| Scale | "Add partitions" | Sizes partitions for 3× peak; knows repartitioning breaks key order | Plans cells, chargeback and a migration path off legacy brokers |
| Ownership | Platform owns the queue | Consumers own handlers and DLQs; platform owns contract, tooling and signals | Writes the standard; decides which use cases the platform refuses |
5.2 Strong Hire Signals#
| Signal | What 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#
| Signal | Why It Misses the Bar |
|---|---|
| "Kafka gives us exactly-once" with no qualification | Confuses a transactional log feature with end-to-end effect on external systems |
| Global ordering by default | Caps throughput and creates a single point of head-of-line blocking without asking whether order is needed |
| No answer for consumer crash mid-processing | The defining failure of a queue was never considered |
| DLQ mentioned without owner, alert or retention | The design has a place where messages silently die |
| Spends 15 minutes on segment files and indexes | Redesigning a solved storage layer instead of the delivery contract |
| Autoscale consumers as the answer to every backlog | Ignores 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.bytesand 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#
| Phase | Time | Goal |
|---|---|---|
| Framing and intents | 0–4 min | Commit to work queue; capture rate, size, processing time, outage tolerance |
| Entities and API | 4–7 min | Message ID, receipt, lease, ack/nack/extend |
| High-level architecture | 7–12 min | Front-end, log, dispatcher, metadata — no more |
| Ack path | 12–20 min | Quorum, min in-sync, producer idempotence |
| Consumer contract | 20–30 min | Leases, heartbeats, idempotent effect, redelivery |
| Poison + DLQ + ordering | 30–37 min | Retry tiers, key-parking, ownership |
| Multi-tenancy / scale | 37–42 min | Quotas, cells, partition sizing |
| Wrap-up | 42–45 min | Contract summary, what's next, what you skipped |
6.2 How Interviewers Pivot — And What They're Testing#
| Pivot | What They're Testing | Strong Response |
|---|---|---|
| "Now make it exactly-once." | Whether you'll overclaim | Define exactly-once effect; show the consumer transaction; note where broker transactions do and don't help |
| "One customer needs strict global ordering." | Pushback with cost | Single-partition queue for that tenant, throughput ceiling in their SLO, or versioned events instead |
| "Consumers are 6 hours behind." | Backlog reasoning | Age vs retention, drain ETA, downstream capacity before adding consumers |
| "A whole AZ goes down." | Durability mechanics | Quorum ack means no loss; leaders move in ~2–10 s; producers retry; consumers see brief redelivery |
| "Make it 10× bigger." | Scale without re-architecture | Cells, partition headroom, replication bandwidth cost |
| "Messages are 50 MB videos." | Payload boundaries | Claim-check: blob storage + pointer; the queue carries metadata |
| "Delay a message by 3 days." | Knowing what queues are bad at | Scheduler 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#
- "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.
- "How do you implement delayed delivery?" — Delay tiers (fixed-delay queues) for retries; a time-indexed store for arbitrary delays beyond ~15 min.
- "How do consumers dedupe without an unbounded table?" — Dedup window bounded by retention + max redelivery horizon; TTL the table at ~2× retention.
- "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.
- "How would you support priority?" — Separate queues per priority with weighted polling in the consumer SDK; not per-message priority inside a partition.
- "How do you bill tenants?" — Stored GB-hours + produced and consumed GB + request count; replication factor multiplies stored cost.
- "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_syncchanges 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:
- Create the new queue with 128 partitions.
- Producers dual-write — or better, flip producers via config to the new queue at a known point.
- 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.
- 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
| Phase | What 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. |
| Triage | Database 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 fix | Producer-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. |
| Guardrails | Alert handler.duration_p99 > 0.5 × lease_seconds. Autoscaler max bound by DB connection budget (max_pods = db_conn_limit / conns_per_pod). |
| Post-mortem | Why 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
| Phase | What to Do |
|---|---|
| Immediate (0–5 min) | Extend DLQ retention from 14 to 21 days. Confirm no new messages expire. |
| Triage | Group by error: 376K TemplateRenderError: missing locale field; 4K transient SMTP 4xx. Deploy 13 days ago added a locale-dependent field. |
| Quick fix | Fix 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. |
| Guardrails | Alert dlq.oldest_message_age > 1h routes to the email team's pager. dlq.depth added to the team's service dashboard. |
| Post-mortem | Why 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_idconstraint 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
| Phase | What 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. |
| Guardrails | Skip 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_syncset 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
| Phase | What to Do |
|---|---|
| Immediate (0–5 min) | Audit all partitions for min_in_sync < 2 and unclean.election = true. Fix any found. |
| Triage | Change 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 fix | Producers with outboxes (orders, payments) replay the window — 1.4M recovered. Remaining 0.5M (analytics, notifications) declared lost with product sign-off. |
| Guardrails | Config 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-mortem | Five 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
| Phase | What 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. |
| Guardrails | mirror.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#
| Claim | Evidence |
|---|---|
| No queue can deliver exactly once across an external side effect | The consumer can crash after the effect and before the ack; the broker cannot know which happened |
| Broker transactions solve a narrower problem | They 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 cheap | A 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#
| Claim | Evidence |
|---|---|
| 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 cost | Throughput per key is serial; poison messages block; repartitioning becomes a migration |
| Ordering is rarely end-to-end anyway | Producer 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#
| Claim | Evidence |
|---|---|
| Broker failures are well handled by mature systems | Leader election and replication are decades-old solved problems in production logs |
| Most queue incidents start in consumers | Lease storms, poison messages, rebalance storms, unowned DLQs and autoscaling into dependencies all originate consumer-side |
| Platform teams invest in the wrong layer | Broker 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#
| Claim | Evidence |
|---|---|
| Most internal job queues are small | Many services enqueue well under 1K jobs/s |
| Transactional enqueue removes a whole failure class | The job commits with the business row; no outbox, no dual write |
SKIP LOCKED makes polling contention manageable | Workers 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#
| Claim | Evidence |
|---|---|
| Messages reach the DLQ exactly when humans are slowest | Failures cluster around deploys, incidents and weekends |
| DLQ storage is cheap | DLQ volume is typically < 0.1% of primary volume |
| Copying the primary's retention is the common default | Which 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.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Every team chooses | Speed; teams pick what they know | 5+ systems, N on-call rotations, inconsistent semantics at team boundaries, no shared tooling | On-call engineers (toil), the org (duplicated infra), customers (seam incidents) |
| One platform for everything | One contract, one SDK, one on-call | Forces ephemeral fan-out and analytics firehoses into a durability model built for jobs; platform becomes a bottleneck | Teams with mismatched workloads (cost, latency); platform team (every request is theirs) |
| Curated set: one log, one work queue, one pub/sub; shared contract | Each intent gets a fitting system; SDK and conventions unify the seams | Requires a standard and enforcement; three systems still to run | Platform 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.
| Scale | Throughput | Infra ($/month) | Headcount | On-call Load | Dominant Cost Driver |
|---|---|---|---|---|---|
| Startup | ~2K msg/s, 4 MB/s | ~$1–3K managed queue/log | 0.5 eng (shared infra team) | Shared rotation, < 1 page/month | Managed 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 team | Dedicated rotation, 3–6 pages/month, most consumer-caused | Cross-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 tenants | Replication 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#
One-Way Doors vs Two-Way Doors#
| Decision | Door Type | Reversibility Cost |
|---|---|---|
| Partition count and key for an ordered queue | One-way (without key→partition indirection) | Migration with drain-and-cutover per queue; per-tenant coordination |
| Message ID and envelope format (headers, schema reference) | One-way | Every producer and consumer changes; dedup tables keyed on the old ID |
| Delivery semantics promised to consumers | One-way | Weakening a guarantee silently breaks consumers that relied on it |
| Exposing broker-native APIs vs a platform API | One-way-ish | Teams bind to broker clients; swapping brokers becomes N migrations |
| Choice of broker behind the platform API | Two-way (if the API abstracts it) | Months of migration, but invisible to tenants |
| Lease length, retry schedule, DLQ thresholds | Two-way | Config change per queue |
| Retention per tier | Two-way upward; one-way downward | Shortening retention deletes data immediately |
| Tiered storage on/off | Two-way | Background 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#
| Signal | What 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 ackedin_flight— map of offset → (lease_expiry, generation, consumer_id)acked_above_commit— sparse set of acked offsets above the commit pointattempts— 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#
| Mechanism | Protects Against | Doesn't Protect Against | Cost |
|---|---|---|---|
Quorum ack (min_in_sync=2) | Leader host/AZ loss after ack | Two-AZ correlated loss before flush | ~5–10 ms latency; 2× cross-AZ bandwidth |
| Producer idempotence (pid + seq) | Duplicate append on retry | Producer restart re-sending | Negligible |
| Broker dedup window (by message_id) | Producer restarts within window | Re-sends after window (e.g., 5 min) | Memory for window of IDs |
| Consumer dedup table | All duplicate effects | Handlers with side effects outside the transaction | One row per message; TTL cleanup |
| Broker transactions (log → log) | Duplicate output in read-process-write pipelines | External side effects | Throughput cost; transaction coordinator |
| Leases with generations | Stale acks from expired consumers | Duplicate effect from the expired consumer's completed work | Generation tracking |
| Outbox in producer | Losing an event between DB commit and publish | — | Relay lag 100 ms–s |
Appendix D: API Contract & Client Behavior#
| Behavior | Rule |
|---|---|
| Produce retries | Exponential 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 Requests | Honor Retry-After; never retry immediately; surface quota exhaustion as a metric |
| Receive | Long-poll up to 20 s; batch up to 10; never tight-loop on empty responses |
| Ack | After side effect commits; batch acks every 100 ms or 100 messages |
| Nack | With retry_after for known transient errors; with permanent=true for validation errors → straight to DLQ |
| Heartbeat | Every lease/3; stop heartbeating if the handler exceeds max lease — let it be redelivered |
| Shutdown | Stop receiving, finish in-flight within grace period, nack the rest with zero delay |
| Thundering herd after outage | Consumers 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):
| Metric | Why It Matters | Alert |
|---|---|---|
queue.oldest_unacked_age (per partition, max) | The only lag metric that maps to user impact and retention risk | Critical: > 5 min; standard: > 30 min |
queue.expired_unprocessed_count | Silent data loss | Any non-zero on critical/standard |
queue.redelivery_rate | Lease misconfiguration, consumer crashes, storms | > 5% for 5 min |
ack.stale_receipt_rate | Leases shorter than real processing | > 0.1% |
dlq.oldest_message_age | Unowned or unnoticed failures | Critical: 15 min; standard: 1 h |
partition.isr_size | Durability below standard | < 2 for 2 min pages |
produce.latency_p99 | Ack-path health | > 2× baseline for 5 min |
tenant.quota_rejections | Tenant at limit | Routed 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#
| Scale | What Works | What Breaks Next |
|---|---|---|
| < 1K msg/s | Postgres SKIP LOCKED table or a managed queue | Vacuum pressure, polling load |
| 1K–50K msg/s | Managed queue for jobs; managed log for events | Managed costs; per-queue throughput caps on FIFO modes |
| 50K–500K msg/s | Self-run log + dispatcher layer, one cell, quotas | Noisy neighbors; single metadata group; cross-AZ bill |
| 500K+ msg/s | Multiple cells, tiered storage, placement service, chargeback | Cross-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.
| Quota | Enforcement Point | On Violation |
|---|---|---|
| Produce bytes/s, requests/s | Front end, local token bucket, synced every 1 s | 429 with Retry-After |
| Consume bytes/s | Front end | Throttled receive (smaller batches) |
| Stored bytes per queue | Storage tier, checked per segment roll | Produce 429 for that queue only |
| Partitions per tenant | Control plane at creation | Creation rejected |
| In-flight leases per queue | Dispatcher | Receive 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.