System Design Interview

Course Content

System Design Interview

31 sections · 71 lessons

Message Queue: requirements, scale and the queue-versus-log choice


"Design a distributed message queue." Ten sections in this course have drawn a queue as a box with an arrow going in and an arrow coming out. This section opens the box.

This lesson covers the first half of the interview: what the prompt actually means, the requirements and the arithmetic, and the one decision the whole design turns on — whether you are building a queue or a log. The next two lessons go deep on partitions and storage, then on delivery guarantees and replication.

Which product is behind the box?Whichproduct is it?Queue or log?Ordering guarantee?Retention window?Delivery semantics?Consumers per topic?
Ten earlier lessons drew this as a box with two arrows; the first question is which of the two products it was.

Two different products share one name

The prompt is ambiguous because the industry uses "message queue" for two things that behave differently.

A task queue holds work items. A consumer takes one, does the work, and the message is deleted. Amazon SQS and RabbitMQ are the common implementations. The mental model is a to-do list that shrinks.

A distributed log is an append-only file that many consumers read independently, each tracking its own position. Kafka is the common implementation; Pulsar and Redpanda are others. The mental model is a shared journal that everybody reads at their own pace and nobody erases.

They are not variations on a theme. They have different storage designs, different failure behaviour, and different things they cannot do. Choosing between them is the first decision of the design, and the second half of this lesson makes the case.

The questions that shape everything after

  1. Queue or log semantics? Is a message gone once consumed, or retained so a second consumer group can read it later, and so a buggy consumer can be replayed after a fix?
  2. What ordering guarantee? Total order across all messages means a single writer and a throughput ceiling. Order within a key — all events for one order ID, in sequence — is usually what the business actually needs.
  3. What delivery guarantee? At-most-once, at-least-once, or exactly-once. The answer changes where offsets get committed and whether consumers need to be idempotent, meaning safe to run twice with the same input.
  4. How long is data retained? "Until consumed" and "seven days regardless" produce storage estimates that differ by orders of magnitude.
  5. How large are messages, and how many producers and consumers? A 1 KB event stream and a 10 MB video-frame stream are different systems.

If there is time, two more: single datacentre or several, and is this multi-tenant with per-team quotas?

The assumptions this section uses

QuestionAssumption
SemanticsLog — retained, replayable, multiple independent consumer groups
OrderingGuaranteed within a partition, not across partitions
DeliveryAt-least-once by default, with a path to effective exactly-once
RetentionSeven days by time, with an optional size cap
Message size1 KB average, 1 MB hard maximum

Requirements and scale

With the questions answered, put numbers down first, because the storage design is the hard part of this system and storage is what the numbers size.

Storage is the hard part, so size it first1 M per secondPartitioned topicsabout 1 KB1 GB/s ingest100 MB/sper brokerTens of brokers7 daysDisk, not memoryabout 600 TBSequential writesNumberConsequenceMessagesMessage sizeThroughputRetentionDisk needed
Seven days of retention at a gigabyte a second is the number that turns this into a storage design.

Functional requirements

  • Producers publish messages to a named topic.
  • Consumers subscribe to a topic and receive messages, in order within a partition.
  • Messages are retained for a configured window whether or not anyone consumed them.
  • Multiple independent consumer groups can read the same topic without affecting each other.
  • A consumer can rewind to an earlier position and reprocess.

Non-functional requirements

  • High throughput for streaming workloads, low latency for task workloads. Both, ideally, and the design should say which it favours.
  • Durable: an acknowledged message survives a broker dying.
  • Available: the cluster keeps accepting writes when a minority of brokers are down.
  • Horizontally scalable: adding brokers adds capacity.

The arithmetic

Assume a target of 100,000 messages per second sustained, averaging 1 KB each. Invented numbers, but the shape is realistic for a mid-sized company's central event bus.

Ingest bandwidth. 100,000 × 1 KB = 100 MB/s written into the cluster.

Daily volume. Round a day to 100,000 seconds instead of 86,400 — a trick from Rounding aggressively and staying fast, which costs about 15% accuracy and saves 30 seconds. 100 MB/s × 100,000 s = 10 TB per day.

Retained volume. Seven days of retention = 70 TB. With a replication factor of three, so each partition has one leader and two followers, that is 210 TB of disk.

Broker count. If each broker carries 12 × 16 TB drives, roughly 190 TB raw and call it 150 TB usable after formatting and headroom, two brokers hold the data. But disk is not the binding constraint here — network and recovery time are. At 210 TB across two machines, one machine failing means re-replicating 105 TB, which at 1 GB/s takes 29 hours. Spread across 20 brokers, each holds about 10 TB, and a rebuild moves 10 TB in under three hours.

Network out. Replication multiplies writes by three: 100 MB/s in becomes 300 MB/s of inter-broker traffic. If three consumer groups each read the full stream, that is another 300 MB/s out. Total cluster traffic ≈ 700 MB/s ≈ 5.6 Gbit/s, or 280 Mbit/s per broker across twenty. Comfortable on 10 Gbit networking.

The consumer side

100,000 messages per second across, say, 100 partitions is 1,000 messages per second per partition. A consumer doing 5 ms of work per message handles 200 per second, so each partition needs five consumer processes — except that a partition can only be read by one consumer in a group at a time. That constraint reappears in Topics, partitions, and ordering, and it is the reason partition count is a capacity-planning decision rather than an afterthought.

Queue versus log

Back to the first question, because it is the distinction the whole design turns on, and the one candidates most often blur.

The queue model

A message arrives, sits in a buffer, is handed to exactly one consumer, and is deleted once that consumer acknowledges it. If five consumers are attached, each message goes to one of them — they compete for work.

This is a good model for tasks. "Resize this image" should happen once. Adding consumers adds throughput with no coordination. The broker's storage stays small because it holds only the backlog.

What it cannot do: replay. Once acknowledged, the message is gone. If the image resizer had a bug for two hours, the inputs are unrecoverable. It also cannot serve two different consumers that both need every message — a billing service and an analytics service reading the same order stream — without the broker duplicating the message into two queues at publish time.

The log model

A message is appended to a file and stays there for the retention window. Each consumer group records an offset: the position of the next message it has not read. Reading does not delete. Two groups reading the same topic have two offsets and never interact.

Queue semanticsa message is consumed and goneconsumeracknowledged and deletedCompeting consumers share the work. Adding a consumer addsthroughput. Nobody can re-read what has been consumed.Log semanticsrecords stay; each reader keeps an offset0123456reader A · offset 2reader B · offset 5Readers are independent and can rewind. Retention is a time or size policy, nota consequence of reading.which one the question wants"Send this email once" is a queue. "Every downstream team needs this event stream, and the analytics team wants to reprocess lastmonth" is a log. Naming which one you are building is half the answer.
A queue destroys on read and a log does not — every other difference follows from that one.

Why the log model won for event streaming

Three reasons, none of them about performance.

Replay is a bug-fix tool. Fix the consumer, reset the offset, reprocess three days of events. In the queue model, that data no longer exists.

Adding a consumer is free. A new team wants the order stream. In the log model they pick a group ID and start reading. In the queue model, someone has to change the publisher or add a fan-out binding.

The storage design is simpler and faster. Append-only files with no per-message state mutation, which Storage on disk shows is what makes 100 MB/s sustainable on ordinary disks.

What the log costs you

Be honest about this, because interviewers probe it.

  • Storage. 70 TB retained versus a queue's few gigabytes of backlog.
  • No per-message acknowledgement. You cannot say "message 47 failed, redeliver only that one". Offsets move forwards, so a poison message blocks its partition until you skip it or route it to a dead-letter topic yourself.
  • Deleting one record is hard. Right-to-erasure requests against an immutable log need compaction with tombstones, or encryption with per-user keys you can destroy.

Recommendation: build the log, and implement task-queue behaviour on top of it when needed. The reverse — retrofitting replay onto a queue — is not possible.