Javathoughts Logo
Javathoughts
Published on
Views

Design a Distributed Queue System Design Interview Guide

Authors
  • avatar
    Name
    Javed Shaikh
    Twitter

← System Design Interview Preparation

This guide walks through Design a Distributed Queue the way you would in a backend or Java interview. The numbers are interview estimates. They help you show your thinking. They are not a production capacity plan.

You are designing something in the family of Kafka / SQS: a log or queue other services use. Product systems in this series already use queues; here you build one.


1. Problem

Producers want to hand off work. Consumers want to pull it later. They should not share a database table and hope.

Example:

  • Payment service publishes PaymentSucceeded
  • Three consumer groups: email, ledger, analytics
  • Messages survive a broker restart
  • If email crashes, it resumes from an offset
  • Poison messages go to a dead-letter queue after N retries

2. Functional Requirements / FR

RequirementWhat it means
ProduceAppend a message to a topic/queue.
ConsumePull batches.
Topics / queuesNamed streams.
PartitionsScale a topic horizontally.
OrderingPer partition (or per key).
At-least-onceDefault delivery.
AcknowledgementsCommit offset / delete after process.
RetriesRedeliver unacked messages.
Dead-letter queueAfter max attempts.
Consumer groupsEach group gets all messages; members split partitions.
Persistence + replicationDisk + followers.

Out of scope: full exactly-once transactions across DBs, SQL on the log.


3. Non-Functional Requirements / NFR

RequirementWhy it matters
High write throughputSequential disk.
Durableacks=majority.
Low produce latencyBatching vs delay.
Independent scaleAdd consumers without slowing producers.
BackpressureSlow consumers lag; producers still write (until disk fills).

Interview line: append-only log, partitions for scale, offsets for consumers, replication for durability.


4. Back-of-the-Envelope Calculation

Traffic assumptions

  • 1 billion messages/day
  • Average payload 1 KB
  • Read/write: each message consumed by 2 groups → ~2× reads

QPS

1,000,000,000 / 86,400 ≈ 11,600 produce QPS average
If peak is 5x, design for around 58,000 QPS.

Consume ≈ 116,000 QPS peak across groups (2×).

Storage

1,000,000,000 × 1 KB ≈ 1 TB/day
7-day retention ≈ 7 TB plus replicas (×3 → ~21 TB)

Cache / memory

Hot tail of the log in page cache. Offset map in memory on the broker.

Servers

Sequential writes: a few brokers can take 58k QPS of 1 KB if batched. Interview: start with 6 brokers (3 AZ × 2), add partitions. Estimates, not exact production numbers.


5. APIs

POST /v1/topics/{topic}/messages
{ "key", "value", "headers" }

GET  /v1/topics/{topic}/partitions/{p}/messages?offset=&limit=

POST /v1/groups/{group}/offsets
{ "topic", "partition", "offset" }

POST /v1/topics  { "name", "partitions", "replicationFactor" }

Ack: produce returns offset. Consume is pull (long poll).


6. Data Model

Commit log (per partition)

Segment files: offset → bytes. Immutable except the active segment.

Metadata

Topic, partition count, leader, replicas. Stored in a metadata quorum (or a small cluster).

Offsets

(group, topic, partition) → offset. Separate from the log so consumers do not rewrite messages.

DLQ topic

Same shape, plus original_topic, attempts, error.


7. High-Level Design

Think log, not a linked list in Redis (Redis lists are a junior answer for this scale).

Distributed Queue architecture

Distributed Queue architectureProducers append to a replicated commit log. Consumers pull by offset. Failed messages after retries go to a dead-letter queue.appendreplicatepullcommit offsetfailed⚙️Producer🌐Queue API / Broker🗄️Commit Log / Partitions🗄️Replica👷Consumers🧠Offsets / Metadata☠️Dead-letter Queue
Producers append to a replicated commit log. Consumers pull by offset. Failed messages after retries go to a dead-letter queue.

Components:

  • Producer → Queue API / Broker
  • Broker appends Commit Log / Partitions
  • Replication to followers
  • Consumers pull
  • Offsets / metadata store
  • Failed processing → DLQ

Produce path: hash key → partition → leader append → wait replica acks → return offset.

Consume path: group coordinator assigns partitions. Client fetches, processes, commits offset.


8. Deep Dives

Ordering

Same key → same partition → order. No global order unless 1 partition (does not scale).

At-least-once

Crash after process before commit → redelivery. Consumers must be idempotent. Point to Pulse-Guard Kafka.

Acknowledgements

  • Produce acks: 1 vs all
  • Consume: auto vs manual commit (prefer manual after success)

Retries + DLQ

Broker redelivers uncommitted. Application retries with backoff; then produce to DLQ and commit the original to avoid a loop.

Consumer groups

N consumers in one group split partitions. You cannot have more useful consumers than partitions.

Persistence

fsync policy: every message vs every 10ms. Interview: batch fsync, replication for durability.


9. Bottlenecks

  • Hot partition (bad key, e.g. null or one user)
  • Small messages without batching
  • Disk full (retention)
  • Rebalance storms when consumers flap
  • Huge payloads (cap 1 MB)

10. Tradeoffs

ChoiceUpsideDownside
Queue (delete on ack)SimpleNo replay
Log (retain)Replay, multi-groupDisk
At-least-onceSimpleDupes
Exactly-once (tx)CleanerHard, slower
Push vs pullPush feels livePull scales better

Pick: replicated log, pull consumers, at-least-once, DLQ, partitions by key.


11. Failure Modes

FailureHandling
Leader diesFollower elected; short unavailability
Slow replicaRemove from ISR; do not wait forever
Consumer poisonDLQ after N
Duplicate produceIdempotent producer (pid + seq) if you have time to mention it
Split brainQuorum metadata, not two leaders

12. Interview Answer in 10 Minutes

"I would design a distributed queue as a partitioned, replicated commit log.

A billion 1 KB messages a day is about 12,000 QPS, about 60,000 at 5× peak, about 1 TB a day, about 7 TB a week before replicas. Producers append to a partition by key. Brokers replicate. Consumers pull and store offsets separately so many consumer groups can read the same data.

Delivery is at-least-once. Handlers are idempotent. After retries, poison goes to a DLQ. Order is per partition.

If a broker dies, a replica becomes leader. If one key is hot, that partition is the bottleneck — choose keys carefully."


13. Interview Talking Points

  • Log vs queue.
  • Partition = scale + order.
  • Offset commits.
  • Replication / acks.
  • Idempotent consumers.
  • Metrics: produce p99, ISR, lag per group, DLQ rate, disk.
  • Event-driven architecture sits on top of this.

14. Follow-up Questions

SQS vs Kafka?
SQS: managed queue, easy, weaker fan-out. Kafka: log, replay, high throughput.

How would Java consumers work?
Spring Kafka, manual ack, idempotency store. See scaling event-driven Java.

Exactly-once?
Idempotent produce + transactional consume-transform-produce. Heavy; mention, do not oversell.

Retention 0 vs 7 days?
Queue semantics vs replay for new consumers.


Related JavaThoughts reading:


Next in this series: Design an API Gateway.