AI System Design
← All chapters
Chapter 19 6 min read

Design a Distributed Message Queue

A distributed message queue decouples producers from consumers and absorbs bursts, but modern designs in the Kafka mould go further: they persist messages in a replicated, append-only log so consumers can replay history. The design rests on partitioned topics, sequential disk IO with batching, consumer groups with rebalancing, and leader-follower replication whose acknowledgement settings trade latency for durability.

Architecture at a glance
  1. Producer (client-side routing + buffer)
  2. Partition leader broker
  3. Write-ahead log segments on disk
  4. Follower replicas (ISR)
  5. Consumer group (pull)
  6. Coordinator + offset state storage
  7. Metadata / ZooKeeper

Problem & requirements

Classic message queues such as RabbitMQ deliver messages and delete them once consumed, while event streaming platforms such as Kafka and Pulsar keep messages for a configurable retention period so multiple consumers can read and re-read them. The line between them has blurred; this design targets the streaming end, because it subsumes the simpler queue.

Requirements: producers send messages, consumers receive them, and messages can be consumed repeatedly or only once. Retention is configurable, say two weeks. Messages are kilobytes in size. Ordering is preserved in the order they were produced to a partition. Delivery semantics (at-most-once, at-least-once, exactly-once) are configurable per use case. The system must offer high throughput or low latency as configured, scale horizontally, and be durable.

Messaging models and core concepts

In point-to-point messaging each message is consumed by exactly one consumer, then disappears. In publish-subscribe, messages are sent to a topic and every subscriber receives a copy. A topic is split into partitions, which are the unit of parallelism and ordering: each partition is an ordered, append-only sequence of messages, each identified by an offset. Partitions are spread across brokers, the servers that store them.

Consumers are organised into consumer groups. Within a group, each partition is assigned to exactly one consumer, which gives point-to-point semantics inside the group and preserves ordering per partition. Different groups each get the full stream, which gives pub/sub. The trade-off is that a group can have at most as many active consumers as there are partitions, so partition count caps parallelism.

  • Messages with the same key go to the same partition, so per-key ordering holds.
  • Global ordering across a topic requires a single partition, sacrificing parallelism.

High-level design

Producers push to brokers; consumers pull from brokers. Brokers hold partitions, each with three kinds of storage. Data storage holds the messages themselves. State storage holds consumer progress, mainly the committed offset per partition per group. Metadata storage holds topic configuration: partition count, retention, replica placement. A coordination service such as ZooKeeper (or Kafka's newer built-in Raft-based KRaft) handles service discovery, controller election and cluster membership.

Keeping these concerns separate matters because they have different access patterns. Data is huge, append-only and read mostly sequentially. State is tiny but updated constantly and must survive consumer restarts. Metadata changes rarely but must be strongly consistent. Each can therefore use the storage technique that fits.

Figure 1One partition, replicated and consumed
One partition, replicated and consumedISR OF ORDERS-P0CONSUMER GROUP Aproducefetchfetch (behind)pull batchcommit offsetcommitelect leaderProducerkey → partition P0Broker 1 · P0 leaderappends, sets offsetsBroker 2 · P0 followercaught upConsumer A1owns P0Broker 3 · P0 followerlagging · out of ISRCoordination servicebrokers, leaders, ISRConsumer A2owns P1 (not shown)State storagecommitted offsets
Producers write only to the partition leader; followers pull from it, and only those that keep up form the ISR. Within a consumer group each partition has exactly one owner, whose committed offset lets a replacement resume where it stopped.

Deep dive: storage with WAL, segments and page cache

The workload is write-heavy and read-heavy, with no updates or deletes, and reads are mostly sequential. A general-purpose database is a poor fit for both patterns at scale. Instead, each partition is a write-ahead log: a plain file to which new messages are appended. Because files cannot grow forever, the log is split into segments; only the newest is active for writes, and old segments are deleted or compacted once past retention.

This exploits a well-known fact: rotational and solid-state disks are slow at random access but very fast at sequential access, often hundreds of megabytes per second. Appends are sequential, and consumers reading near the tail are served from the operating system's page cache, so most reads never touch the disk. Data on disk is stored in the same format it travels over the wire, avoiding copy and conversion; zero-copy transfer from file to socket further cuts CPU.

Deep dive: message format, batching and the producer flow

A message carries a key (used to choose the partition, not unique), a value payload, the topic and partition, an offset assigned by the broker, a timestamp, a size and a CRC for integrity. Keeping the format identical end to end means brokers can move bytes without parsing them.

Batching is pervasive and is the main throughput lever. Producers group messages in memory per partition and send them together, brokers write batches to the log in one call, and consumers fetch batches. Larger batches mean fewer network round trips and larger sequential writes, but each message waits longer. The producer embeds the routing layer in the client library: it fetches partition-to-leader metadata, picks the partition from the key, buffers messages, and sends batches directly to the leader broker, avoiding an extra network hop and proxy.

Figure 2Message record on disk
Message record on diskOffset8 bytesset by brokerTimestamp8 bytesms since epochSize4 bytesrecord lengthCRC4 bytesintegrity checkKey8 bytesvariable · picks partitionValue32 bytesvariable · payloadtotal 64 bytes
Brokers store and ship records in exactly the format producers send, so bytes move from network to log to consumer without re-encoding. The CRC lets any hop detect corruption; key and value are variable-length (typical sizes shown).

Deep dive: consumers, rebalancing and replication

Consumers pull rather than being pushed to. Pull lets each consumer control its own rate, so a slow consumer is not overwhelmed, and makes aggressive batching natural; its downside, busy polling on an empty partition, is solved with long polling. Each group has a coordinator broker that receives heartbeats. When a consumer joins, leaves or misses heartbeats, the coordinator triggers a rebalance: one consumer is elected group leader, computes a new partition assignment, and the coordinator distributes it. Committed offsets live in state storage, so a reassigned partition resumes where the previous owner stopped.

For durability every partition has replicas on different brokers. Producers write only to the leader; followers pull from it. The in-sync replicas (ISR) are those caught up within a configured lag. The producer's acks setting chooses the trade-off: acks=0 never waits and can lose data; acks=1 waits for the leader and loses data if it dies before followers copy it; acks=all waits for every ISR member, giving the strongest durability at the highest latency.

Figure 3Produce with acks=all
Produce with acks=allProducerBroker 1 (leader)Broker 2Broker 31. buffer + compress batch for P02. ProduceRequest(P0, batch, acks=all)3. append to log: offsets 42–494. fetch from offset 425. records 42–496. fetch from offset 427. records 42–498. fetch from 50 (confirms up to 49)9. fetch from 50 (confirms up to 49)10. all ISR caught up: high watermark = 5011. ack: base offset 42
The leader acknowledges only after every in-sync follower has fetched the batch, which advances the high watermark; consumers never see records above it. A slow follower is dropped from the ISR rather than stalling every write.

Scalability and delivery semantics

Producers scale freely because they hold no coordination state. Consumers scale by adding group members up to the partition count, with rebalancing moving work. Brokers scale by adding nodes and moving partition replicas onto them; a safe approach temporarily adds new replicas, lets them catch up, then removes old ones. Partition counts can grow, but doing so changes which partition a key maps to, breaking per-key ordering for existing data, so pick generous counts up front.

Delivery semantics follow from when acknowledgements and offset commits happen. At-most-once: producer does not retry and consumer commits before processing; fine for metrics where loss is acceptable. At-least-once: producer retries with acks=1 or all and consumer commits after processing, so duplicates are possible and consumers should be idempotent. Exactly-once requires idempotent producers with sequence numbers plus transactional commits of output and offsets; it is the most expensive and suits payments and accounting.

Advanced features, failure handling & wrap-up

Consumers often care about only some messages. Rather than creating a topic per subtype, brokers can support filtering by tags in message metadata, so the broker returns only matching messages without decoding payloads. Delayed or scheduled messages, such as cancelling an unpaid order after 30 minutes, are first written to internal temporary topics and moved to the real topic when due, using timing wheels or delay-level topics to avoid scanning.

On failure, a dead leader is replaced by an ISR member chosen by the controller, so committed data survives; a dead consumer's partitions are rebalanced to peers; and a slow follower falls out of the ISR so it cannot stall writes. The overall lesson is that a log on disk, partitioned and replicated, with simple clients that batch, gives both throughput and durability without a database.

Key numbers

Default retention example
2 weeks
Sequential disk throughput
hundreds of MB/s
Max consumers per group
= number of partitions
Ack levels
acks = 0 / 1 / all
Typical message size
~KB

Key terms

Topic
A named stream of messages to which producers write and consumers subscribe.
Partition
An ordered, append-only shard of a topic that is the unit of parallelism and ordering.
Offset
A message's monotonically increasing position within its partition.
Consumer group
A set of consumers that share a topic's partitions so each message is processed once per group.
Rebalance
Reassignment of partitions among group members when membership changes.
In-sync replicas (ISR)
Replicas sufficiently caught up with the leader to be eligible for commit and leader election.
Segment
A bounded file of a partition's log; old segments are deleted when past retention.
Exactly-once
A delivery guarantee achieved with idempotent producers and transactional commits of results and offsets.

Common mistakes

  • Using a general-purpose database as message storage and losing the benefit of sequential IO.
  • Adding more consumers than partitions and expecting more throughput.
  • Increasing partition count on a keyed topic without realising it breaks per-key ordering.
  • Claiming exactly-once delivery without idempotent producers and atomic offset commits.
  • Using acks=1 for financial data and losing messages on leader failure.
  • Push-based delivery that overwhelms slow consumers.

Further study

  • Kafka: a Distributed Messaging System for Log Processing (2011)
  • The Log: What every software engineer should know about real-time data's unifying abstraction (Jay Kreps)
  • Apache Pulsar
  • RocketMQ
  • KIP-500: Replace ZooKeeper with a Self-Managed Metadata Quorum

Now practise it