Design a Distributed Message Queue
A Kafka-style log that ingests a million messages a second, keeps them in order per key, replicates them so no acknowledged message is lost, and lets many consumer groups read at their own pace.
Last updated 2026-09-30. Difficulty: hard. Patterns: append-only-log, partitioning, replication, consumer-groups, exactly-once. Reported at LinkedIn, Amazon, Confluent, Uber, Microsoft.
The interviewer asks, the candidate answers and draws, and you press Next. Pause to answer yourself at the key decisions, and ask the AI Mentor anything along the way.
Functional requirements
- Produce messages to a topic. Each message has an optional key, a value and headers. Messages with the same key stay in order.
- Consume in groups. A consumer group shares a topic’s partitions among its members; each message is processed by one member of each group. Many independent groups can read the same topic.
- Replay. Consumers can rewind to an earlier offset or timestamp and reprocess, for example after a bug fix.
- Retention. Messages are kept for a configured time or size (say 7 days), or compacted to the latest value per key, regardless of whether they were consumed.
- Delivery guarantees. At least once by default; exactly-once processing for consume-transform-produce pipelines.
- Out of scope. Per-message routing rules and priorities of a classic broker (RabbitMQ), delayed messages, and cross-datacentre mirroring beyond a mention.
Non-functional requirements
- Throughput (1 M msgs/s · ~1 GB/s in). Spread over hundreds of partitions; a single partition handles roughly 10 MB/s.
- Durability (no acknowledged message lost with one broker down). This is the requirement that drives the hardest trade-off: acknowledging only after several replicas have a message costs latency, and choosing which replica may become leader after a failure decides whether data can be lost.
- Latency (p99 < 20 ms produce-to-consume). Within a region, for streaming pipelines. Batch consumers do not care.
- Ordering (per partition (per key)). No global order: that would cap throughput at one machine.
- Availability (tolerate a broker failure with seconds of disruption). Leaders fail over automatically; producers and consumers retry transparently.
- Fan-out (many consumer groups per topic). Reading must not slow writing, and slow consumers must not affect fast ones.
Back-of-envelope estimates
- Ingress: ~1 GB/s. 1 M messages/s × ~1 KB = 1 GB/s produced, before replication.
- Disk writes with replication: ~3 GB/s. Replication factor 3: every byte is written by the leader and two followers, 3 GB/s of sequential writes across the cluster.
- Retention storage: ~1.8 PB. 1 GB/s × 604 800 s (7 days) ≈ 605 TB × 3 replicas ≈ 1.8 PB. Storage, not throughput, decides the broker count, which is why old segments go to object storage.
- Minimum partitions: ~100. A partition sustains roughly 10 MB/s of writes comfortably: 1 GB/s ÷ 10 MB/s = 100 partitions. In practice several hundred, so consumer groups have enough parallelism.
- Network out per broker: ~300 MB/s. Total traffic: 1 GB/s in + 2 GB/s replication + 3 GB/s to three consumer groups = 6 GB/s. Across 20 brokers that is 300 MB/s each, comfortable on 25 Gbit/s NICs.
- Messages per batch: ~100. Producers batch for a few milliseconds: at 10 k msgs/s per producer, a 10 ms linger gives 100 messages per request, turning a million tiny writes into thousands of large ones.
Components
- Producers (batching · idempotent): Services writing events. The client library hashes the key to choose a partition, batches messages per partition for a few milliseconds, compresses, and sends to that partition’s leader with a sequence number for deduplication.
- Consumer group: Instances of one application sharing a topic. Each owns some partitions, fetches from its committed offset, processes, and commits progress. Adding instances (up to the partition count) adds parallelism.
- Schema registry: Stores message schemas (Avro, Protobuf) by id and checks compatibility, so producers can evolve formats without breaking consumers. Messages carry a schema id, not the schema.
- Controller quorum (KRaft / Raft): Holds cluster metadata (topics, partitions, which broker leads each, the in-sync replica sets) in a replicated log, detects broker failures and elects new partition leaders.
- Group coordinator: A broker role that manages one consumer group’s membership: tracks heartbeats, triggers rebalances, assigns partitions, and stores committed offsets.
- Partition leader (broker 1): The broker currently leading a partition. Appends producer batches to the log, serves follower fetches and consumer fetches, and advances the high watermark once all in-sync replicas have a message.
- Log segments (page cache → disk): Each partition is a directory of append-only segment files (~1 GB) with sparse offset and timestamp indexes. Writes go to the OS page cache and are flushed in the background; reads of recent data come straight from memory.
- Committed offsets (__consumer_offsets topic): Each group’s progress per partition, stored as messages in an internal compacted topic, so it is replicated and durable like any other data.
- Follower (broker 2): Fetches from the leader continuously, like a consumer, and appends the same bytes at the same offsets. Stays in the in-sync replica set while it keeps up.
- Follower (broker 3): The third replica, in a different rack or zone. With acks=all and min.insync.replicas=2, a write survives the loss of any one broker.
- Tiered storage (S3): Closed segments older than a few hours are copied to object storage and deleted locally. Consumers reading far back fetch from it transparently. Cuts broker disk needs by an order of magnitude.
- Log cleaner (retention · compaction): Deletes segments past their retention time or size, and for compacted topics rewrites old segments keeping only the latest message per key (and tombstones for a while).
User flows
- Produce a message with acks=all. The write path: batch on the client, append sequentially on the leader, acknowledge once every in-sync replica has it.
- The producer fetches cluster metadata and caches which broker leads each partition. Clients talk to partition leaders directly; there is no load balancer in the data path. Metadata is refreshed on errors such as "not leader".
- It serialises the message with a registered schema and picks the partition from the key. partition = hash(key) mod partitions, so all events for one order or one user land in one partition, in order. Messages without a key are spread in sticky batches.
- Messages for each partition are batched for a few milliseconds, compressed, and sent to the leader. Batching turns a million small messages a second into thousands of large requests. The idempotent producer attaches its producer id and a per-partition sequence number so the broker can discard duplicates on retry.
- The leader appends the batch to the partition’s active segment. An append to the end of a file: sequential I/O that lands in the OS page cache. The broker does not fsync per message; durability comes from replication across machines instead.
- Followers fetch the new data and append it at the same offsets. Replication is pull-based: followers send fetch requests continuously, which also tells the leader how far each has got.
- When every in-sync replica has the batch, the high watermark advances and the producer is acknowledged. With min.insync.replicas = 2, the write is refused if fewer than two replicas are in sync, rather than accepted with a single copy. Consumers only ever see messages below the high watermark.
- Consume as a group. Partitions are divided among the group’s members; each reads sequentially from its committed offset and commits progress after processing.
- Consumers join the group; the coordinator assigns partitions among them. Each partition goes to exactly one member, which is what preserves per-key order. With cooperative rebalancing, members only give up the partitions that move.
- Each consumer fetches from its committed offset, up to the high watermark. Recent data is served from the page cache and sent with zero-copy (sendfile), straight from the file to the socket. A consumer that is far behind reads older segments from disk or tiered storage without disturbing others.
- The consumer processes the batch, then commits the next offset. Committing after processing gives at-least-once: a crash between processing and commit means those messages are processed again, so processing must be idempotent. Committing before processing would give at-most-once.
- If a member dies, its heartbeats stop and its partitions are reassigned. The new owner starts from the last committed offset, so anything processed but not committed is processed again. Lag per partition (high watermark minus committed offset) is the key health metric.
- One key floods a partition. The scale-breaking case: per-key ordering means a single hot key cannot be spread, so one partition, one leader and one consumer take all of it.
- A single tenant starts producing 50 MB/s, all with the same key. Its partition receives five times what a partition comfortably sustains; the other 299 partitions are fine. The leader’s disk and network for that partition saturate.
- Followers fall behind on that partition and drop out of the in-sync set. When the ISR shrinks below min.insync.replicas, acks=all writes to that partition are rejected: the system protects durability by refusing writes instead of accepting them unreplicated.
- The one consumer that owns the partition builds lag. Adding consumers does not help: one partition is read by one member of the group. Lag on that partition grows while the rest of the group idles.
- The fix is in the key: split the hot key into sub-keys, or throttle the tenant. If the tenant’s events only need order per entity within the tenant, key by tenant and entity. Per-client quotas on the brokers stop one producer from starving others in the meantime.
- A partition leader’s broker dies. The failure path: the controller notices, promotes an in-sync follower, and clients retry. With acks=all and no unclean election, nothing acknowledged is lost.
- Broker 1 stops heartbeating to the controller quorum. After the session timeout (a few seconds) the broker is declared dead. Every partition it led needs a new leader.
- The controller picks a new leader from each partition’s in-sync replicas and publishes the change. Only an in-sync replica is eligible, so the new leader has every acknowledged message. Electing an out-of-sync replica ("unclean election") would restore availability faster but lose data, so it is disabled.
- The new leader’s log is the truth; any replica ahead of it truncates to the high watermark. Messages the old leader had written but not yet replicated to all in-sync replicas were never acknowledged, so dropping them breaks no promise. Leader epochs make truncation exact.
- Producers and consumers get "not leader" errors, refresh metadata, and retry against the new leader. The idempotent producer’s sequence numbers make retries safe: if a batch had in fact been written, the new leader recognises the sequence and acknowledges without duplicating it.
- Retention, compaction and tiered storage. The background path. Data is kept by time or size, not by consumption, and old segments move to cheap storage.
- When the active segment reaches ~1 GB or an hour, it is closed and a new one starts. Only the active segment is written to. Closed segments are immutable, which makes deleting, copying and compacting them simple.
- Closed segments older than a few hours are copied to object storage and removed from local disk. Brokers keep only recent data locally; 1.8 PB of retention mostly lives in S3. Consumers replaying last week fetch from it transparently, a little slower.
- The log cleaner deletes segments past retention, or compacts keyed topics. Deletion is per segment, so it is just removing files. Compaction rewrites old segments keeping only the newest message per key, turning a changelog into a snapshot of current state that never grows beyond the number of keys.
Deep dives
Why a log
Classic queues delete a message once it is acknowledged. Why keep an append-only log instead?
A traditional broker tracks every message’s delivery state and deletes it on acknowledgement. That means random writes and per-message bookkeeping, and once a message is consumed it is gone: a second application cannot read it, and a bug fix cannot replay it.
An append-only log flips this: the broker only appends and the consumers keep track of their own position (an offset). Writes and reads are sequential, which disks and the OS page cache handle at hundreds of megabytes a second, and any number of consumer groups can read the same data independently.
- Partitioned append-only log; consumers track offsets; retention by time or size chosen
- Classic broker with per-message acks (RabbitMQ, SQS) situational: task queues where each job is processed once and order does not matter
- Store messages in a database table and poll rejected
The answer: Each topic is split into partitions; each partition is an append-only sequence of segment files with sparse offset and timestamp indexes. Producers append, consumers read sequentially from an offset they choose and commit, and data expires by retention, not by consumption. Writes rely on the OS page cache rather than per-message fsync, with durability coming from replication; reads of recent data are served from memory with zero-copy transfer from file to socket. This is the design LinkedIn published as Kafka.
Why not fsync every message?
An fsync per write would cap throughput at the disk’s sync rate, a few thousand a second. Instead, a message is acknowledged once it is in the page cache of several brokers in different racks. Losing it would need all of them to lose power before flushing, which replication across failure domains makes very unlikely.
What is zero-copy and why does it matter?
The sendfile system call lets the kernel send file data directly from the page cache to the network socket, without copying it into the broker’s memory and back. With many consumer groups reading the same recent data, it cuts CPU and memory traffic dramatically.
How would you support delayed messages or per-message retries?
Not in the log itself. Common patterns: a retry topic per delay (retry-1m, retry-10m) that a consumer re-publishes from after waiting, and a dead-letter topic for messages that keep failing. Or use a classic queue for that workload.
How do you find the message at a given offset quickly?
Each segment is named by its first offset, so a binary search over segment names finds the file. A sparse index inside the segment maps every few kilobytes of offsets to a byte position, so the read starts near the right place and scans a little.
Partitions and ordering
How do you choose partitions, and what ordering can you promise?
One partition lives on one leader and is consumed by one member of each group, so it caps both write and read throughput (roughly 10 MB/s write, and one consumer’s processing speed). Throughput comes from many partitions in parallel.
Global ordering would require a single sequence, which is a single machine. What applications almost always need is order per entity: all events for one order, one account or one device in sequence. Keying by that entity gives exactly that.
- Order per partition; partition by hash of the message key; over-provision partitions chosen
- Single partition for global order situational: low-volume streams that truly need a total order, such as a ledger’s command log
- Round-robin without keys situational: independent events such as logs and metrics
The answer: Messages carry the entity id as key; the partition is hash(key) mod N. Start with more partitions than today’s throughput needs (several hundred for a 1 GB/s topic) so consumer groups have room to scale and repartitioning is rare, but not tens of thousands, because each partition costs file handles, memory and leader elections. Ordering is promised per key only. When a topic must grow, create a new topic with more partitions and migrate producers and consumers, rather than adding partitions in place to a keyed topic, which would send some keys to new partitions mid-stream.
Why does adding partitions break ordering for keyed topics?
hash(key) mod N changes when N changes, so a key’s new messages go to a different partition while its old ones are still being consumed from the old one. Two consumers may then process the same key concurrently and out of order. Hence: plan partitions ahead, or migrate to a new topic.
A consumer’s processing is slow and per-message. How do you get more parallelism than partitions?
Within one consumer, process different keys in parallel while keeping each key’s messages sequential (a pool keyed by message key), and commit offsets only up to the lowest fully processed message. This decouples processing parallelism from the partition count.
How many partitions is too many?
Each partition has files, index memory and a leader to elect on failure, so tens of thousands per cluster slow down failover and metadata operations, though newer controllers handle far more. A few thousand per cluster and tens to a few hundred per topic is a common, safe range.
Replication without losing data
When exactly is a message safe, and what happens to messages when a leader dies?
Each partition has a leader and followers. If the producer is acknowledged as soon as the leader has the message, a leader crash before followers fetch it loses the message. If the producer waits for every replica, one slow follower stalls all writes.
The in-sync replica set (ISR) is the compromise: the replicas currently keeping up. Acknowledge when all of them have it, drop followers that fall behind, and require a minimum ISR size so durability never silently degrades to one copy.
- acks=all to the ISR, min.insync.replicas = 2 with RF 3, only in-sync replicas may become leader chosen
- acks=1 (leader only) situational: metrics, logs and other data where small loss is acceptable
- Allow unclean leader election situational: availability-first streams where a gap is better than downtime
The answer: Replication factor 3, replicas in three racks or zones. Producers use acks=all and the broker requires min.insync.replicas = 2, so an acknowledged message is on at least two brokers, and writes fail loudly rather than proceed with one copy. Followers fetch continuously; one that falls behind for longer than a configured time leaves the ISR. The high watermark (the offset every ISR member has) bounds what consumers see. On leader failure, the controller elects a new leader from the ISR only, and followers truncate any unacknowledged tail using leader epochs. Unclean election stays off for important topics.
Why do consumers only see messages below the high watermark?
A message above it might not survive a leader failure: if the leader dies, the new leader may not have it, and it would vanish. A consumer that had already read it would have processed data that no longer exists. Limiting reads to the high watermark prevents this.
Two brokers in a three-replica partition die. What happens?
The ISR shrinks to one, below min.insync.replicas, so acks=all writes to that partition are rejected until a replica recovers. Reads of already-committed data continue from the survivor. That is the deliberate choice: unavailable for writes rather than accepting unreplicated data.
How do you replicate across regions?
Not with synchronous replicas: cross-region latency would sit on every write. Use asynchronous mirroring (MirrorMaker or cluster linking) into a separate cluster, accept seconds of lag, and design consumers for failover with offset translation.
Consumer groups and offsets
How do several instances of a service share a topic, and how do they know where they left off?
Each message must be processed once per application, but an application runs as many instances that come and go. Assigning partitions to instances gives each a disjoint share and preserves per-key order, but every time an instance joins or leaves, partitions must move.
Progress is just an offset per partition. Where it is stored, and when it is committed relative to processing, decides the delivery guarantee.
- Coordinator-managed groups, cooperative sticky assignment, offsets in a replicated internal topic, commit after processing chosen
- Offsets stored by the consumer in its own database, in the same transaction as its results situational: consumers that write results to their own database
- Stop-the-world rebalancing (eager) rejected
The answer: Each group has a coordinator broker. Members heartbeat; a missing member or a new one triggers a rebalance, and the cooperative sticky assignor keeps existing assignments and moves only what must move, in incremental steps. Committed offsets are written to the compacted __consumer_offsets topic, replicated like any data. Consumers commit after processing a batch, giving at-least-once delivery, and make processing idempotent (upserts keyed by message id, or deduplication tables). Lag per partition, the gap between the high watermark and the committed offset, is the main health metric for every consumer.
A consumer takes 10 minutes on one message and gets kicked out of the group. Why?
If it does not poll within max.poll.interval, the coordinator assumes it is stuck and reassigns its partitions; another member then reprocesses the same message, and the cycle repeats. Fix: hand long work to a separate pool and keep polling (pausing the partition), or raise the interval for that workload.
How do you replay a day of events after fixing a bug?
Reset the group’s offsets to a timestamp (the broker’s time index finds the matching offset per partition), with the consumers stopped, then restart them. Because processing is idempotent, reprocessing is safe. Retention must cover the window you want to replay.
What does consumer lag tell you?
Lag that grows steadily means consumers cannot keep up and need more instances or faster processing. Lag on one partition only points to a hot key or a poison message. Lag measured in time (how old the next unprocessed message is) is easier to alert on than in messages.
Exactly-once
Producers retry and consumers crash. How can a pipeline process each message exactly once?
Duplicates come from two places. A producer retrying after a lost acknowledgement may write the same message twice. A consumer that processes a message and crashes before committing its offset will process it again after restart.
True exactly-once delivery over a network is impossible, but exactly-once effects are achievable: deduplicate on write, and make "process, write results, commit offset" atomic.
- Idempotent producer + transactions for consume-transform-produce; read_committed consumers chosen
- At-least-once plus idempotent consumers situational: results written to external systems (databases, APIs) that support upserts or idempotency keys
- At-most-once (commit before processing) situational: telemetry where loss is preferable to duplicates
The answer: Producers are idempotent: each has a producer id and a sequence number per partition, and the leader rejects a batch whose sequence it has already written. For pipelines that read from one topic and write to others, a transaction groups the output writes and the input offset commit; a transaction coordinator writes commit or abort markers, and consumers in read_committed mode never see aborted or in-flight messages. For effects outside the log (a database, an email), use at-least-once and make the effect idempotent with the message id as an idempotency key, or store offsets in the same database transaction as the result.
Does exactly-once in the queue stop a consumer sending two emails?
No. The queue can only make its own writes atomic. If processing sends an email and then crashes before committing, the email goes out again. Make the email service idempotent (key by message id) or record sent ids; the queue cannot cover side effects it does not own.
What does the idempotent producer cost?
Very little: a producer id, a sequence number per batch, and the broker remembering the last few sequences per producer per partition. It is enabled by default in modern clients because it removes a whole class of duplicates for free.
A transactional producer crashes mid-transaction. What do consumers see?
Nothing from that transaction: read_committed consumers only read up to the last stable offset, below any open transaction. When the producer restarts with the same transactional id, the coordinator aborts the old transaction (fencing the zombie by epoch), and the aborted messages are skipped.
How do you know exactly-once is actually working?
Count duplicates downstream on a sample keyed by message id, and reconcile totals end to end (orders in, invoices out) per hour. A pipeline that claims exactly-once but drifts in reconciliation has a side effect outside the transaction.
Storing a week of data
1.8 PB of retention with replication. How do you store it without buying racks of disks?
Retention is what makes replay and late consumers possible, but at 1 GB/s a week of data with three replicas is 1.8 PB. Keeping it all on broker disks means sizing the cluster for storage, not throughput, and makes every broker replacement copy tens of terabytes.
Almost all reads are of recent data (consumers keeping up). Old data is read rarely, for replays and backfills, where a little extra latency is fine.
- Tiered storage: recent segments on local disk, older ones in object storage; deletion and compaction per segment chosen
- Everything on local disks situational: short retention (hours) or small volumes
- Compaction only situational: changelog topics where only current state matters
The answer: Segments roll at ~1 GB or one hour. Closed segments older than a few hours are uploaded to object storage and deleted locally, so each broker holds roughly a day of its partitions’ data on disk while the full week lives in S3. Retention deletes whole segments, which is cheap. Changelog topics use compaction, keeping the newest message per key and tombstones for a grace period. Consumers reading old offsets fetch from object storage through the broker transparently.
Does tiered storage weaken durability?
No: segments are only uploaded once closed and replicated, and object storage is more durable than a broker disk. The local replicas still provide durability for recent data, which is where acknowledged writes live.
How does compaction handle deletes?
A delete is a message with the key and a null value (a tombstone). Compaction keeps the tombstone long enough for consumers to see it, then removes both the tombstone and the earlier values, so the key disappears from the snapshot.
Why roll segments by time as well as size?
A low-volume partition might take weeks to fill a gigabyte, so its oldest messages could never be deleted or tiered on schedule. Rolling at least hourly makes retention and tiering apply on time regardless of volume.
How long does a broker replacement take with tiered storage?
Only the recent local data (around a day of its partitions) must be copied from other replicas, instead of a full week. That cuts recovery from many hours to well under an hour and reduces the window of reduced redundancy.
The theory behind it
- Message queues and event streams: When to go asynchronous, queues vs logs (SQS, RabbitMQ vs Kafka), delivery guarantees, ordering, consumer groups, dead-letter queues, backpressure, and the outbox pattern.
- Replication and consistency: Leader-follower, multi-leader and leaderless replication; synchronous vs asynchronous; quorums; CAP and PACELC; consistency models from eventual to linearizable; read-your-writes and failover.
- Distributed transactions and idempotency: Why cross-service transactions are hard, two-phase commit vs sagas, the outbox pattern, idempotency keys, exactly-once semantics, and how to keep money and inventory correct.