Design a distributed message queue: durable logs, consumer ownership and replay
Clarify work queue versus replayable stream, then design a partitioned log with independent consumer groups. Trace acknowledgment loss, offset commits, leader failover and poison records without claiming arbitrary external effects are exactly once.
1. Resolve queue versus stream before choosing storage
When you say message queue, do you want a task to disappear after one worker handles it, or several applications to independently replay the same retained events? I will ask that first because deletion, acknowledgment and storage are different. For this chapter we agree on a replayable partitioned log.
| Candidate asks | Illustrative interviewer reply | Design consequence |
|---|---|---|
| Who consumes a message? | Several independent applications; one active consumer per partition within a group. | Persist data separately from each group’s progress. |
| What ordering is required? | Order for one entity key, not across the entire system. | Stable key-to-partition mapping and sequential processing for that key. |
| How long can consumers replay? | Seven days. | Retention is time-based; a slow consumer can fall behind the retained range. |
| What is the workload? | 100,000 messages/s average, 1 million/s peak, 1 KB payload. | Size storage from the average; network and burst handling from peak. |
| What durability is needed? | Acknowledged records survive one broker or zone failure. | Replicate before acknowledgment; refuse writes when the required safe replicas are unavailable. |
| Do we need exactly-once payments? | No; explain how external consumers avoid duplicate effects. | Stable event IDs and transactional consumer dedup; log guarantees stop at defined boundaries. |
| Are global transactions required now? | No; discuss them as an extension. | Baseline has per-partition commits and group progress, not an invented all-system transaction. |
2. Define publish and consume guarantees
My publish acknowledgment means the record is committed under the chosen replication policy. It does not mean any consumer has processed it. Consumers commit progress only after their effect is durably handled, which permits replay after a crash.
| Requirement | Agreed target or scope | Design implication |
|---|---|---|
| Publish | Batch append, key routing, bounded record size and optional producer idempotency. | Validate before append; distinguish rejected, committed and unknown outcomes. |
| Consume | Pull batches, independent groups, replay from retained offsets. | Broker backpressure follows available bytes; clients control processing concurrency. |
| Ordering | Append order within a partition; no global order. | Parallel work inside one partition needs ordered commits or per-key sequencing. |
| Latency | Illustrative p99 publish target 20 ms in-region under admitted load. | Replication and batching consume the budget; do not claim the same latency across continents. |
| Durability | Survive one replica/zone loss; disallow unsafe leader promotion. | Unavailable partitions may stop writes rather than acknowledge data that can be lost. |
| Failure processing | Bounded retries and explicit poison-record policy. | Skipping to a DLQ changes ordering semantics and must be agreed with the consumer. |
3. Separate sustained volume from peak throughput
I use decimal KB and TB. A one-million-per-second peak is not a one-million-per-second seven-day average. I would validate the peak duration and consumer catch-up demand before choosing broker count.
| Quantity | Calculation | Interpretation |
|---|---|---|
| Average ingress | 100,000/s × 1 KB = 100 MB/s; 8.64 TB/day. | Seven days = 60.48 TB raw payload. |
| Replicated retention | 60.48 TB × 3 copies = 181.44 TB. | With an assumed 2:1 compression ratio: 90.72 TB before indexes, headroom and temporary rebalance copies. |
| Peak ingress | 1M/s × 1 KB = 1 GB/s. | Two follower copies add 2 GB/s replication payload; disks and NICs need burst headroom. |
| Two full-rate consumer groups | 2 × 100 MB/s = 200 MB/s average egress. | At peak both reading current traffic need 2 GB/s, before replay and catch-up. |
| Partitions | 256 logical partitions: about 391 messages/s average and 3,906/s peak each if uniform. | An entity with disproportionate traffic still creates a hot partition. |
| One-hour stalled consumer | 100,000/s × 3,600 = 360M records, 360 GB payload backlog. | At 200,000/s processing while 100,000/s new traffic arrives, catch-up takes another hour. |
4. Model logs, ownership and committed positions
A partition has one current leader epoch. Offsets identify records within that partition, not globally. Consumer group generations protect ownership changes; they do not by themselves prevent a stale process from writing to an unrelated external database.
| Record and key | Fields / invariant | Access pattern |
|---|---|---|
| Topic(topicId) | partition count, retention, schema rules, quotas and access policy. | Administrative control plane, versioned changes. |
| Partition(topicId, partitionId) | leader, epoch, replica set and committed boundary. | Metadata lookup and fenced leader changes. |
| Record(partitionId, offset) | eventId, key, timestamp, schemaVersion, headers, payload checksum. | Sequential append and offset-based fetch. |
| ProducerState(producerId, epoch, partition) | last accepted sequence and recent result metadata. | Detect duplicate transport retries within the supported session/window. |
| GroupOffset(groupId, partitionId) | next offset to process and ownership generation. | Commit only contiguous completed progress under current generation. |
| ConsumerDedup(eventId) | Application-owned result or processed marker. | Atomically commit with the business effect where the sink supports transactions. |
5. Define protocol-level failure and replay behavior
These are logical RPC contracts, not a claim of wire compatibility with Kafka. The client refreshes partition metadata on a leader-epoch error, but retries the same producer identity and sequence when supported so a transport retry does not become an intentional new event.
| Operation | Contract | Failure or retry behavior |
|---|---|---|
| Produce(topic, partition, producerEpoch, sequence, records) | Return committed offset range and current leader epoch after safe replication. | On timeout outcome is unknown; retry with the same identity. Reject stale epochs and oversized batches. |
| Fetch(partition, offset, maxBytes, waitMs) | Return committed records and retained range. | OFFSET_EXPIRED requires explicit reset or archive recovery, never silent success. |
| JoinGroup(groupId, memberId) | Return partition assignment and generation. | Heartbeats renew membership; stale generations cannot commit group offsets. |
| CommitOffset(group, partition, nextOffset, generation) | Persist contiguous completed progress. | Commit is idempotent for the same position; do not advance past incomplete earlier work. |
| Replay(group or newGroup, offset/time) | Authorized reset within retention. | Pause current ownership or use a separate group to avoid racing active processing. |
6. Draw the data plane and the metadata authority
I use sequential partition logs replicated across three zones and a separate consensus-backed metadata controller. The controller elects leaders and assigns epochs; it does not route every payload. Consumers pull from committed logs and keep their own progress per group.
7. Trace publication through an uncertain acknowledgment
The producer routes an entity key to a partition, batches records and sends them to the current leader. The leader validates its epoch, sequence and schema, appends records and replicates them before returning committed offsets. A durable acknowledgment policy must define flush and replication behavior, not just count machines.
Route consistently
Keep a stable mapping for entity ordering. Adding partitions can move a key; migrate through an explicit epoch/drain protocol or accept that order spans two partitions.
Append and replicate
Write a checksummed segment and replicate to the required safe set. Advance the committed boundary only when the policy is satisfied; followers cannot expose uncommitted tail records as durable.
Retry without inventing new intent
A lost acknowledgment is ambiguous. Producer sequence dedup covers supported retries, while a stable business eventId lets consumers recognize duplicates after producer restarts or outside that dedup window.
Fence leader replacement
A new leader is selected from safe replicas with a new epoch. An old leader must reject writes once fenced; if safe replication is unavailable, stop rather than promote stale data silently.
Bound disk growth
Enforce producer quotas and retention. If disk pressure threatens safe operation, throttle or reject new writes with an explicit error; do not acknowledge and discard accepted records.
8. Trace processing, rebalancing and poison records
A consumer receives an assignment, reads from its committed next offset, performs work and advances only contiguous completed progress. If it processes offsets 42 and 43 concurrently but 42 is unfinished, committing 44 can lose 42 after a crash. That is why concurrency and commit order must be designed together.
Apply business idempotency
When possible, write the effect and unique eventId marker in one sink transaction. A check-then-effect-then-marker sequence is not atomic and can duplicate effects.
Handle external APIs
If the sink is an external payment or email API, use its idempotency contract or an outbox/adapter and reconciliation. A broker transaction cannot atomically commit an arbitrary remote side effect.
Rebalance safely
Stop fetching revoked partitions, drain bounded in-flight work and commit safe progress. Fence stale group commits; use sink-side dedup or fencing because an old consumer may still finish work.
Choose poison policy
Strict ordered consumers pause the partition and alert after bounded retries. A skip-to-DLQ policy records the original coordinates and changes the contract; persist quarantine before advancing past it.
Recover expired offsets
If retention deletes needed records, surface the gap. Rebuild from an authoritative snapshot/archive with a checkpoint, rather than resetting to latest and claiming no data loss.
9. Discuss throughput without hiding ordering costs
Sequential I/O, batches, compression and page cache help throughput, but storage and network still have finite budgets. I separate producer quotas, live-consumer capacity and replay capacity so a backfill cannot starve current traffic.
| Decision | Chosen baseline | Alternative and trade-off |
|---|---|---|
| Partition count | Enough parallelism for measured throughput and consumer concurrency. | Too many partitions increase metadata, file handles and recovery work; one hot key still serializes. |
| Batches | Bounded batch bytes and linger time. | Larger batches improve amortization and compression but add latency and retry size. |
| Work queue alternative | For independent tasks, per-message visibility and acknowledgment can fit better. | A retained log simplifies replay and fan-out but requires per-group offset and poison-ordering decisions. |
| Exactly-once extension | Transactions can couple reads and writes inside a supported transactional ecosystem. | They do not automatically include arbitrary databases, webhooks or the recipient’s behavior. |
| Cross-region | Asynchronous disaster-recovery copy with a stated lag/RPO. | Synchronous inter-region commits improve recovery guarantees at WAN latency and availability cost. |
10. Walk failures from producer to sink
I observe under-replicated partitions, oldest uncommitted append, leader changes, publish p99, consumer lag in seconds and bytes, retention margin, disk headroom and duplicate-effect prevention. Lag count without message size or processing cost can be misleading.
| Failure | Detection | Recovery and remaining limitation |
|---|---|---|
| Leader dies after append | Epoch change and uncertain producer response. | Elect a safe replica and retry; uncommitted tail may be discarded, committed data must remain under the declared fault model. |
| Insufficient safe replicas | Replication health fails the write policy. | Reject new acknowledgments; availability is sacrificed rather than asserting durability without replicas. |
| Consumer crashes after effect | Offset still points to completed event. | Replay; sink transaction recognizes eventId and avoids repeating the business effect. |
| Stale consumer after rebalance | Old group generation attempts commit. | Reject offset commit; prevent stale sink effects through idempotency/fencing. |
| Poison record | Repeated deterministic failure at one offset. | Pause or quarantine under the agreed policy; DLQ is not a universal ordering-preserving fix. |
| Disk corruption | Checksum failure or replica mismatch. | Repair from a verified replica, isolate bad storage and test restore; replication is not a substitute for backup against logical corruption. |
| Retention overtakes consumer | Requested offset precedes earliest retained offset. | Report data gap and restore/rebuild explicitly; never silently claim successful delivery. |
11. Finish at the end-to-end correctness boundary
The broker durably stores an ordered partition log; consumer groups independently track progress; applications make replay safe. I can now explain what happens at every lost-ack boundary instead of using exactly once as a blanket label.
I would validate the design with crash injection around append acknowledgment, effect commit and offset commit, plus a retention-overrun drill. Capacity benchmarks must include failover, skew and replay, not just a uniform producer benchmark.
Technical references
Primary references explain underlying mechanisms. Workloads and architecture choices above remain proposed interview assumptions.