Design a Distributed Message Queue
Interviewers ask this to see whether you understand the systems you use every day. The question is "design something like Kafka or a managed queue service": a durable, scalable broker that producers write to and consumers read from. It tests partitioning, replication, the storage format of a log, delivery guarantees, consumer coordination and failure. The Kafka chapter explains how to use such a system. This chapter explains how to build one.
It follows the usual shape: understand the problem, set up the interface, build the high-level design, then go deep on the questions interviewers use to separate levels.
1. Understanding the problem
A producer sends messages to a named stream (a topic). One or more consumers read the messages. The broker stores messages durably in between, so producers and consumers can run at different speeds and fail independently.
Functional requirements
Core:
- Producers publish messages to a topic.
- Consumers read messages from a topic, in order within a partition, at their own pace.
- Messages are stored durably until a retention limit, and are not lost if a broker fails.
Confirm in or out: multiple independent consumer groups, replay of old messages, message ordering guarantees, delayed delivery, priorities, per-message acknowledgement versus offset commits, and message size limits. A sensible opening: "I will design a partitioned, replicated log with consumer groups and offset tracking, with at-least-once delivery, and treat delays and priorities as extensions."
Non-functional requirements
- High throughput. Hundreds of thousands to millions of messages per second across the cluster.
- Low latency. Milliseconds from publish to availability.
- Durability. An acknowledged message survives the loss of a broker.
- Availability. Producers and consumers keep working through broker failures.
- Scalability. Add brokers to increase capacity.
- Ordering within a partition.
Estimation
Assume 1 million messages per second at peak, 1 KB each, retained for seven days, replicated three times.
| Quantity | Calculation | Result |
|---|---|---|
| Ingress | KB | about 1 GB per second |
| Daily volume | 1 GB/s 86,400 | about 86 TB per day |
| Stored, 7 days, 3 copies | TB | about 1.8 PB |
| Per broker, 50 brokers | about 36 TB each | |
| Partitions | say 1,000 across the cluster | 1,000 messages per second per partition |
What the numbers say. Gigabytes per second and petabytes of retention mean the data path must be sequential disk and network, not random access, and the system must be spread across many brokers. Per-partition rates are modest, which is why partitioning is the scaling mechanism. State these as assumptions to adjust.
2. The set up
Core entities
- Topic with a number of partitions.
- Partition: an ordered, append-only log of messages.
- Message: key, value, timestamp, headers, and an offset within its partition.
- Broker: a server that stores some partitions.
- Consumer group: a set of consumers sharing the work of a topic, with a committed offset per partition.
API
produce(topic, key, value, acks) -> { partition, offset }
fetch(topic, partition, offset, max) -> [ messages ]
commit(group, topic, partition, offset) -> ok
join_group(group, topic) -> assigned partitions
3. High-level design
- Brokers store partitions as logs on disk and serve produce and fetch requests.
- A metadata and controller service, backed by a consensus protocol, keeps the cluster's state: which brokers exist, which partitions each hosts, who the leader of each partition is, and the topic configuration. It handles leader election when a broker fails.
- Producers discover the partition leaders from metadata and send messages directly to the leader of the chosen partition.
- Consumers in a group are assigned partitions, fetch from the leaders and commit their positions to an offset store.
4. Potential deep dives
Deep dive 1: How do you store messages so writes and reads are fast?
The challenge. Gigabytes per second of writes and many concurrent readers, on ordinary disks.
Weak: store each message as a row in a database, indexed by id. A database does random writes, maintains indexes and locks, and cannot sustain this rate. Deleting expired messages is expensive.
Solid: an append-only log file per partition. Always write at the end of a file, which is a sequential write, the fastest disk operation. A reader supplies an offset and reads forward sequentially.
Excellent: segmented logs with a sparse index, batching and the page cache. Split each partition's log into segment files, named by the offset of their first message. Only the newest segment is written. Keep a sparse index per segment mapping some offsets to byte positions, so a read finds the segment by comparing offsets to file names, then uses the index to jump near the right position and scans forward.
<!--fig:segments-->Why this works well:
- Retention is cheap. Delete whole old segments, with no per-message work.
- Sequential I/O on both write and read.
- Batching. Producers send batches, brokers write them as one unit, and consumers fetch many messages per request, amortising overhead and enabling compression.
- The operating system's page cache serves recent data from memory, so consumers that keep up read from RAM, not disk, and the broker does not need its own large cache.
- Zero-copy transfer sends data from the file cache to the network without copying through application memory.
Deep dive 2: How do you replicate for durability and availability?
The challenge. A broker will fail. Acknowledged messages must survive, and the partition must stay available.
Weak: a single copy per partition. A broker loss loses its partitions and halts them.
Solid: leader and follower replicas. Each partition has a replication factor, such as 3. One replica is the leader, handling all reads and writes, and the others are followers that pull new messages from the leader. If the leader fails, a follower is elected.
Excellent: in-sync replicas, high watermark and tunable acknowledgement. Track the in-sync replica set: the followers that are caught up within a time or lag limit. Define the high watermark as the offset up to which all in-sync replicas have the data. Consumers can read only up to the high watermark, so they never see a message that could be lost in a leader failure.
Let producers choose the acknowledgement level:
| Level | Meaning | Durability | Latency |
|---|---|---|---|
| None | Do not wait | A crash can lose messages | Lowest |
| Leader | Wait for the leader to write | Lost if the leader fails before replication | Low |
| All in-sync | Wait for all in-sync replicas | No acknowledged loss while one in-sync replica survives | Highest |
Pair "all in-sync" with a minimum in-sync replicas setting, such as 2 of 3, so the leader rejects writes when too few replicas are available instead of acknowledging data that is not safely replicated. On leader failure, the controller elects a new leader from the in-sync set, which guarantees no acknowledged message is lost. Electing an out-of-sync replica restores availability sooner at the risk of losing data, and should be off for important topics. Say that this is the availability versus durability trade-off in a concrete form.
Deep dive 3: How do producers choose partitions, and what does that mean for ordering?
The challenge. Ordering costs parallelism.
Weak: one partition for the whole topic. Total ordering, and throughput limited to one broker.
Solid: partition by key. Hash the message key to a partition, so all messages with the same key land in the same partition and stay in order. With no key, spread messages across partitions for balance.
Excellent: a deliberate partitioning strategy. Choose a key with even distribution and meaning for ordering, such as an order or user identifier. Watch for skew, since a hot key makes a hot partition. Choose the partition count from the target throughput and consumer parallelism, with room to grow, because increasing it changes the key-to-partition mapping and breaks ordering for existing keys. State clearly: order is guaranteed only within a partition, never across partitions.
Handle producer retries safely: a retry after a timeout can write a duplicate. Use an idempotent producer with a producer identifier and per-partition sequence numbers, so the broker discards a message it has already stored.
Deep dive 4: Consumer groups and offsets
The challenge. Many consumers share a topic, and each message must be processed by one consumer of a group, in order per partition, and survive consumer failure.
Weak: the broker tracks delivery state per message. Per-message state for millions of messages per second is heavy, and it makes ordering and replay awkward.
Solid: consumers track an offset per partition. A consumer reads from an offset and periodically commits its position. After a restart it resumes from the committed offset. The unit of progress is a single number per partition per group, which is cheap.
Excellent: group coordination, rebalancing and delivery semantics. A group coordinator assigns partitions to the consumers of a group, so that each partition has exactly one consumer in the group, and reassigns them (a rebalance) when consumers join, leave or fail, detected by heartbeats. Partition count caps group parallelism. Delivery semantics follow from when the offset is committed:
- Commit before processing: at-most-once, and a crash loses a message.
- Commit after processing: at-least-once, and a crash causes reprocessing.
- Exactly-once effects: at-least-once plus idempotent processing, or transactions that commit results and offsets atomically.
Store the committed offsets durably and replicated, often in an internal topic itself. Allow consumers to seek to any offset, which gives replay. Reduce the disruption of rebalances by keeping membership stable and by incremental assignment.
Deep dive 5: Failure handling
Broker failure. Leaders on that broker are re-elected from in-sync followers, clients refresh their metadata and retry, and the controller creates new replicas on other brokers to restore the replication factor, with throttling so recovery does not starve live traffic.
Controller or metadata failure. The metadata service runs on a consensus group of an odd number of nodes, so it survives minority failures, and clients continue to use cached metadata for existing partitions.
Network partitions. A leader cut off from the controller must stop accepting writes once its in-sync set shrinks below the minimum, to avoid divergent logs. The new leader's epoch (a leader epoch number) lets followers detect and truncate stale data from the old leader after it returns.
Slow consumers. They lag but do not affect others, because they only read. Monitor lag. If lag exceeds retention, a consumer loses data because messages are deleted before it reads them, so alert before then.
Disk failure. Replicas on other brokers hold the data. Repair by re-replicating.
Poison messages. The broker is agnostic. Consumers implement retry limits and dead-letter topics.
Deep dive 6: Retention, compaction and storage cost
Retention. Delete segments older than a time or beyond a size limit. This is configured per topic.
Log compaction. For topics that represent the latest state per key, keep only the newest message for each key and drop older ones, so a consumer reading from the start sees the current state of every key. A tombstone (a key with a null value) removes a key eventually.
Tiered storage. Keep recent data on local disks and move older segments to cheaper object storage, while still serving fetches for them, so retention can be long without a huge local disk bill.
Compression. Compress batches, which saves network and disk, at the cost of CPU.
Deep dive 7: Delays, priorities and per-message acknowledgement
A pure log does not support them. When a product needs delayed delivery, priorities, per-message acknowledgement with redelivery or competing consumers with individual retry, a different design fits better: a queue with per-message state and a visibility timeout, where a delivered message is hidden until acknowledged or the timeout expires. Explain the trade-off: per-message state costs throughput and complicates ordering, and in return gives individual retries and delays. A single system can offer both models, but they are different mechanisms, and you should say which one the requirements need.
5. What is expected at each level
Mid-level. You describe producers, brokers and consumers, store messages durably and propose partitions for scale and replication for safety.
Senior. You explain the segmented append-only log and why it is fast, leader and follower replication with in-sync replicas, acknowledgement levels, consumer groups with offsets and partition assignment, and delivery semantics. You size the cluster.
Staff. You reason about the high watermark and leader epochs, the controller's consensus group, rebalance costs, compaction and tiered storage, multi-tenant quotas, cross-region replication, and operational concerns such as throttled recovery and capacity planning.
6. Interview questions and model answers
Q: Why is a log fast? Writes are sequential appends and reads are sequential scans, which disks and the page cache handle very efficiently. Batching, compression and zero-copy transfer reduce overhead, and retention is done by deleting whole segment files.
Q: How do you guarantee a message is not lost? Replicate each partition, require acknowledgement from all in-sync replicas with a minimum of two, and elect new leaders only from the in-sync set. Consumers read only up to the high watermark.
Q: How does ordering work? Within a partition, by offset. Messages with the same key go to the same partition. There is no ordering across partitions.
Q: How do consumers share work? A consumer group divides the partitions among its consumers, so each partition has one consumer in the group. Each group tracks its own committed offset per partition, so many groups can read the same topic independently.
Q: What happens if a consumer crashes? The coordinator detects missed heartbeats and reassigns its partitions. The new consumer resumes from the last committed offset, so messages after it are processed again, which is why handlers must be idempotent.
Q: How would you add delayed messages? A pure log does not support them, so I would add a separate delay mechanism that holds messages until due and then publishes them, or use a queue model with per-message state and timers.
7. Common mistakes
- Storing messages in a general-purpose database.
- A single copy of each partition.
- Acknowledging writes before they are replicated, with no minimum in-sync count.
- Claiming ordering across partitions.
- Tracking delivery state per message in a log-style system.
- Electing a leader from out-of-sync replicas by default.
- Forgetting that consumer lag beyond retention means lost data.