Vibe Engines
YouTube
System Design

Design a Message Queue

Step 1 / 9

Learn system design by building a distributed message queue like Kafka step by step.

The numbers to beatappend-onlywrite pathsequentialdisk I/Ooffsetmessage id

The whole design, in writing

Learn system design by building a distributed message queue like Kafka step by step. An interactive guide covering the append-only log, partitions, consumer groups and offsets, replication, leader election, retention, and delivery guarantees.

Every step of the build above, written out: the problem each piece solves, the option that was taken and the ones that were not, the numbers, and how it fails in production.

The big idea

What is a message queue?

One service produces events much faster than another can handle them — or you want many services to react to the same event without the producer knowing they exist. Calling each consumer directly is brittle: a slow or dead consumer drags the producer down with it.

Producerspublish events
New in this step: Producers.

Put a durable log in the middle. Producers append; consumers read at their own speed. The producer fires and forgets; consumers come and go, replay history, and never block the producer. This one primitive underpins event-driven systems everywhere.

What the new pieces do

Producersproducer
Services that emit a firehose of events — clicks, orders, metrics — and want to hand them off and move on, without waiting for anyone to read them.

Step 1 · The core

An append-only log

A producer needs to hand off a message reliably, and a consumer needs to read messages in order, possibly long after they were written. A delete-on-read queue forgets too soon and can’t support replay or multiple readers.

Brokerappend + serveCommit Logappend-only
New in this step: Broker, Commit Log.

Consumers read in order long after a write, and multiple readers want the same stream. Model the topic as…

  1. Delete-on-read forgets too soon: no replay, and a second consumer can’t read what the first consumed. You lose history and multi-reader support.

  2. Random inserts/deletes mean index churn and locking, and you’d hand-roll ordering and cursors. A purpose-built log is far faster and simpler for an ordered stream.

  3. Writes only append (sequential disk I/O — shockingly fast), nothing is deleted on read, and order is implicit in the offset. Any number of consumers read it independently.

Model the topic as an append-only commit log: writes only ever append to the end and get a monotonically increasing offset. Nothing is deleted on read. Sequential disk appends are shockingly fast, and order is implicit in the offset.

Why this piece earns its place

The speed argument quietly assumes consumers read near the tail. An append is sequential, and a recent read is served out of the operating system’s page cache, so the disk mostly only ever sees writes. A consumer that falls far behind — or one that rewinds to replay, which is the feature we just sold — reads cold offsets, and those reads come off the disk in the middle of the file while producers are still appending to the end. Now one machine is doing a sequential write and a scattered read at the same time, and the tail latency every healthy consumer sees gets worse because of the one that isn’t. The defense is to make that visible and bounded: split the log into segments so a cold read touches a bounded file instead of wandering, monitor how far behind each reader is rather than only whether it is alive, and treat a big backfill or a replay from the beginning as scheduled work — throttled, or run off-peak — rather than letting it land silently on the broker that is serving everyone else.

  • append-onlywrite path
  • sequentialdisk I/O
  • offsetmessage id

What the new pieces do

Brokerbackend
The server that accepts writes from producers and serves reads to consumers. In a cluster, many brokers split the work; here we start with one.
Commit Loglog
The heart of it: an immutable, append-only sequence of messages on disk. Writes are just appends; each message gets a monotonically increasing offset.

Back of the envelope

append = sequential disk write
far faster than random I/O — near memory speeds
offset = position in the log
order is implicit, no extra index needed
reads don’t mutate
N consumers read one log independently

Step 2 · Beyond one disk

Partitions for parallelism

A single log is capped by one machine’s disk and CPU — both for write throughput and for how fast consumers can read. One ordered log can’t scale horizontally.

Brokerappend + servePartitionerkey → partitionPartition Logsappend-only
New in this step: Partitioner.

One ordered log is capped by a single machine’s disk and CPU. How do you scale throughput?

  1. Vertical scaling hits a ceiling and is still one machine — one disk’s write rate, one CPU’s read rate. You need to spread a topic across machines.

  2. Each partition is its own log on a possibly different broker; hashing a key keeps related events ordered together while unrelated events spread out. Throughput scales with partition count.

  3. Concurrent writers to one shared log destroy the cheap sequential-append property and reintroduce coordination. Order and speed both come from each partition having a single writer.

Split the topic into partitions, each its own log on a possibly different broker. A Partitioner routes each message — usually by hashing a key — so the same key always lands in the same partition. Throughput scales with partition count.

Why this piece earns its place

The key hash is computed against the current partition count, which quietly makes that count part of the data model rather than a capacity setting. The day you add partitions to keep up with traffic, the same key starts hashing somewhere new: its old events sit in the old partition, its new ones arrive in another, and two consumers each hold half of one entity’s history. Nothing errors — ordering per key simply stops being true across the seam, and anything stateful downstream (a running balance, a last-write-wins projection) is wrong for exactly the keys that moved. So I size partitions with headroom for the parallelism I expect to need, and if I truly must grow, I treat it as a topic migration — produce into a new topic with the new count, let consumers drain the old one to its end, then cut over — not an in-place resize. The other edge is the single hot key: one key is one partition, so a celebrity user is capped at one partition’s throughput however many you add.

  • Npartitions = parallelism
  • key-hashrouting
  • per-partitionordering

What the new pieces do

Partitionerservice
Decides which partition an event lands in, usually by hashing a key. Same key → same partition → ordered together.

Back of the envelope

throughput ∝ partition count
add partitions to add parallelism
same key ⇒ same partition
related events stay ordered
order is per-partition only
not a global order across the topic

Step 3 · Many readers, no deletes

Consumer groups & offsets

If the log never deletes, how does a consumer know what it has and hasn’t read? And how do ten instances of a service share the work without each processing every message?

Consumer GroupBrokerPartition LogsOffsets
New in this step: Consumer Group, Offsets. · swipe to pan the diagram

The log never deletes. How do consumers track progress and share work without each reading everything?

  1. That’s back to a delete-on-read queue — no replay, and only one consumer can ever see a message. The log’s whole value is keeping data and letting readers track their own position.

  2. Re-scanning the whole log per consumer is hugely wasteful with no shared progress. You need a stored cursor, not a full re-read.

  3. Partitions are divided among a group’s members (processed once per group), and a server-side offset records position — so a crash resumes exactly where it left off and a rewind replays. Different groups read the same log independently.

Consumers join a group; the group’s partitions are divided among its members so each message is processed once per group. Each group tracks its position with an offset — a cursor stored by the broker. Different groups read the same log at totally different positions.

Why this piece earns its place

Group membership is the moving part this buys. Every time a member joins, leaves, or stops sending heartbeats, partitions are reassigned — and in the simple protocol the whole group stops consuming while that happens, so a rolling deploy of ten instances can pause a topic ten times rather than once. The dial is the liveness timeout. Set it short and an ordinary long pause — a garbage collection, a slow downstream call inside a batch — evicts a member and triggers a rebalance nobody needed; set it long and a genuinely dead instance keeps its partitions dark for that whole window. I’d set it from the real tail of processing time for one batch rather than from a default, and keep per-batch work bounded so the timeout can stay small. The other half is commit cadence: the offset is committed after the work, so whatever isn’t committed when a member dies is reprocessed by whoever picks the partition up. Commit rarely and that replay is large; commit on every message and you spend the throughput the group was for.

  • per-groupcursor
  • 1×delivery / group
  • rewindto replay

What the new pieces do

Consumer Groupconsumer
Services that read the stream at their own pace — a fraud checker, a search indexer, a data warehouse — each independently, each possibly slow.
Offsetsstore
Where each consumer group has read up to. Stored server-side so a crashed consumer resumes exactly where it left off — the log itself is never mutated.

Step 4 · Don’t lose data

Replication

Partitions live on disks, and disks (and whole brokers) die. If a partition exists on only one broker, its death means permanent data loss and downtime for everyone reading or writing it.

Partition Logsappend-onlyReplicas (ISR)leader + followers
New in this step: Replicas (ISR).

A partition lives on a disk, and disks die. How do you not lose committed data?

  1. Backups are minutes stale, so a broker death still loses recent writes and downs that partition. Durability must be synchronous with the write, not a periodic snapshot.

  2. The leader takes reads/writes; in-sync replicas mirror it, and a write commits only once enough replicas have it (acks=all). Losing the leader loses nothing — an ISR is promoted.

  3. Replicating every partition to every broker wastes enormous disk and write bandwidth. A few copies (typically 3) survive failures — you don’t need a copy on all of them.

Replicate each partition across several brokers. One is the leader (takes all reads/writes); the others are in-sync replicas that mirror it. A write is “committed” once enough replicas have it, so losing the leader loses nothing.

Why this piece earns its place

The guarantee isn’t acks=all on its own — it is acks=all plus a floor on how small the in-sync set is allowed to get. A replica that falls behind is dropped from the ISR, and if that floor is one, the set can shrink to the leader alone while the producer config still says acks=all: every write is then acknowledged by a single machine, the durability you believe you have is gone, and nothing in the write path complains. With ×3 replicas I want the floor at two, so the partition refuses writes instead of pretending — a producer error you can see and retry beats a silent downgrade. Say the cost out loud, because it is real: that setting turns the second broker failure into unavailability for writes on that partition, where the loose setting would have kept taking them. Committed data survives two failures either way; what you are choosing is whether the writes arriving during the outage are honestly rejected or quietly under-replicated.

  • ×3typical replication
  • acks=alldurable write
  • 0committed loss

What the new pieces do

Replicas (ISR)store
Each partition is copied to several brokers. One leader takes writes; in-sync followers mirror them, so a dead broker never loses committed data.

Back of the envelope

×3 replicas typical
survives 2 broker failures
acks=all ⇒ commit on ISR
durability you can trust, slight latency cost
0 committed loss
a promoted ISR already holds every committed write

Step 5 · When a broker dies

Leader election

A leader broker crashes mid-flight. Someone has to notice, promote a healthy in-sync replica to leader, and tell every producer and consumer where to go now — without two brokers both thinking they’re the leader.

Broker Clustermany brokersPartition Logsappend-onlyControllerleader election
New in this step: Controller.

The leader broker for a partition crashes mid-flight. How do clients keep working?

  1. Clients have no consistent global view and could split-brain — two brokers each believing they lead. Leadership must be decided by one authority, not negotiated by clients.

  2. The controller watches health and promotes an in-sync replica via consensus (ZooKeeper/KRaft), then propagates metadata so clients transparently reconnect. Consensus is needed only for who-leads, not every message.

  3. Freezing the partition until a human fixes a machine destroys availability. The point of replication is to fail over in seconds, not wait for repair.

A Controller watches broker health and runs leader election over a consensus store (ZooKeeper, or Kafka’s own KRaft). It promotes an in-sync replica and propagates the new metadata, so clients transparently reconnect to the new leader.

Why this piece earns its place

Election carries one setting that can undo the previous step. When every in-sync replica for a partition is down and only a stale follower is left, the controller has two options: wait, leaving the partition unavailable until an in-sync broker returns, or promote the stale replica and take the partition back online now. That second option is the only path by which committed writes are actually lost in this design — the promoted replica truncates to what it had, and messages consumers were already told existed stop existing. Worse, those offsets get reused, so a returning consumer can read a different message at a position it has already processed. I leave the promotion off by default, because a topic feeding a ledger, a billing job, or anything reconciled against another system can’t afford a hole it cannot detect. I’d enable it deliberately, per topic, for streams where a gap is cheaper than a stall — metrics, click tracking — and only after saying in the room that this is the knob that trades the durability promise for uptime.

What the new pieces do

Controllerservice
Tracks which brokers are alive and elects a new partition leader when one dies. Backed by a consensus store (ZooKeeper, or KRaft today).

Step 6 · The log can’t grow forever

Retention & compaction

An append-only log that never forgets eventually fills every disk. But you also can’t blindly delete — some consumers replay history, and some topics need the latest value per key kept indefinitely.

Partitionerkey → partitionPartition Logsappend-onlyReplicas (ISR)leader + followersRetentiontime / size
New in this step: Retention.

Apply a retention policy: drop log segments older than N days or beyond a size cap. For keyed topics, use log compaction to keep only the most recent message per key, so the log becomes a compact snapshot of current state.

Why this piece earns its place

Compaction has a sharp edge that time-based retention doesn’t. Deleting a key from a compacted topic is itself a record — a null value, a tombstone — and a consumer rebuilding state only learns the key is gone by reading it. But the tombstone can’t live forever either, or compaction would never reclaim anything, so it survives for a bounded window and is then removed along with the key’s history. Any consumer that was offline, or bootstrapping from the start of the log, across that window sees the key’s older values and never the delete — and rebuilds its state holding a record the source deleted months ago. It comes back silently, in one consumer, and reads like a bug in that service. So the window isn’t a disk setting: it is the longest outage or slowest full rebuild you intend to tolerate, and it should be set from that number. If a rebuild would take longer than that, the honest answer is to restore from a snapshot and consume forward, not to replay a compacted log from the beginning.

What the new pieces do

Retentionpolicy
The log can’t grow forever. Old segments are dropped by age or size — or compacted to keep only the latest value per key.

Step 7 · What can it promise?

Delivery guarantees

Networks drop acks and clients crash mid-batch, so producers retry and consumers reprocess — producing duplicates. Applications need to know what the system actually guarantees.

ProducersConsumer GroupBroker ClusterPartitionerPartition LogsReplicas (ISR)OffsetsControllerRetention
The system as it stands at this step. · swipe to pan the diagram

Retries and crashes cause reprocessing. What guarantee do you offer, and how do you avoid duplicates?

  1. Not retrying gives at-most-once — dropped messages on any failure. You can’t get reliability by giving up delivery; exactly-once is built on top of at-least-once, not instead of it.

  2. Committing first means a crash after the commit but before the work silently loses the message (at-most-once). The safe default is to commit after processing.

  3. Retry until acked (accepting dupes), then dedup by producer-id + sequence and use transactions to commit write+offset atomically. Consumers commit offsets only after processing — that ordering defines the guarantee.

Default to at-least-once: retry until acknowledged, accept possible duplicates. Layer idempotent producers (dedup by producer id + sequence) and transactions (atomic write-plus-offset-commit) to reach exactly-once within Kafka. Consumers commit offsets only after successfully processing.

Why this piece earns its place

Exactly-once inside Kafka is a throughput and latency setting, not a switch. Transactions make the write and the offset commit atomic, but consumers have to be configured to read committed data only, which means a record is invisible to them until the producer commits the transaction it belongs to. So the commit interval becomes the end-to-end latency floor for everything downstream: commit per message and you pay a round trip plus a marker in the log for every record; commit in large batches and you get the throughput back while every consumer waits a batch to see anything at all. I’d pick that interval from the slowest downstream consumer’s actual latency budget, and default to plain at-least-once wherever the work is naturally idempotent — an upsert keyed on the event id costs nothing and needs none of this. Worth stating too: producer dedup covers a producer’s own retries, not an application that crashes and regenerates the message afterwards. That duplicate has to be caught by a business key.

You did it

You just designed a message queue.

ProducersConsumer GroupBroker ClusterPartitionerPartition LogsReplicas (ISR)OffsetsControllerRetention
The finished design, end to end. · swipe to pan the diagram

Everything you assembled, in order

  • An append-only commit log — ordered, immutable, replayable.
  • Partitions split the log across brokers; a key-hash keeps order per partition.
  • Consumer groups + server-side offsets: read once per group, rewind to replay.
  • Replication with a leader + in-sync replicas so no committed write is lost.
  • A controller runs leader election over a consensus store on broker failure.
  • Retention and compaction keep the log from growing without bound.
  • At-least-once by default; idempotent producers + transactions for exactly-once.

Where an interviewer pokes next

Getting the boxes right is the easy half. These are the questions that separate a candidate who drew the diagram from one who has run the thing. Answer each one out loud before you open it.

  1. Why is Kafka so fast despite writing to disk?

    Sequential appends + the OS page cache + zero-copy (sendfile) reads. Appending to the end of a file is sequential I/O, which modern disks handle at near-memory throughput, and consumers are often served straight from page cache without copying through user space. The log structure unlocks all of it.

  2. How many partitions should a topic have?

    Enough for your peak parallelism (a partition is the unit of both consumer concurrency and ordering), but not so many that controller metadata, open files and end-to-end latency balloon. You can add partitions later, but it breaks key→partition stability — so size with headroom up front.

  3. What’s the trade-off in choosing a partition key?

    The key controls both ordering and load balance. Too coarse (e.g. country) creates hot partitions; too fine loses the ordering you wanted. Pick the entity whose events must stay ordered (user id, order id) and whose cardinality spreads load evenly.

  4. Does exactly-once work end-to-end, including external sinks?

    Within Kafka, yes — idempotent producers + transactions give exactly-once for read-process-write between topics. To an external system (a database, an email) it’s only exactly-once if that sink is idempotent or joins the transaction; otherwise you fall back to at-least-once + idempotent writes downstream.

  5. A consumer group is slow — what happens to the cluster?

    Nothing upstream: producers keep appending and other groups read normally. The slow group’s lag (latest − committed offset) grows; if it can’t catch up before retention drops old segments, it loses those messages. You monitor consumer lag and scale the group out, up to the partition count.

Check yourself — the answers, and why

Eight steps in, these are the calls you should be able to make cold. Pick one, then read why.

  1. The core data structure of a queue like Kafka is…

    Ordered, immutable, replayable — reads don’t mutate it, so many consumers share it.

  2. Ordering in Kafka is guaranteed…

    Partitions parallelize; a good key keeps related events in one partition, hence ordered.

  3. A crashed consumer resumes correctly because…

    Reading is just advancing an offset; restart resumes from the last committed position.

  4. "acks=all" means a write is acknowledged…

    Committing only after replication means losing the leader loses no committed data.

  5. Kafka’s default delivery guarantee is…

    Retry-until-acked yields duplicates; idempotent producers + transactions layer exactly-once on top.

How you’d open this design in an interview

Before any boxes: agree what it must do, pin the qualities that shape everything, then build — naming each trade-off as you make it. The walkthrough above is that exact order.

What it must do

Agree on these before drawing a single box.

  • Publish + subscribe: producers append events; consumers read the stream at their own pace, replaying history.
  • Partition: split a topic into partitions routed by key hash — same key, same partition, ordered.
  • Track progress: consumer groups split partitions; a server-side offset lets a crash resume and a rewind replay.
  • Durability: replicate each partition (leader + in-sync replicas); commit only after replication.
  • Survive + bound: a controller elects a new leader on failure; retention/compaction cap the log; at-least-once by default.

The qualities that shape everything

Each one names the mechanism that buys it.

Decouple a fast producer from slow or dead consumers
A durable append-only log in the middle — producers append and move on; consumers read at their own speed, and a dead consumer just stops advancing its offset.
Scale throughput past one machine
Split the topic into partitions on different brokers, routed by key hash so same-key events stay ordered while unrelated events spread.
Read once per group, and rewind to replay
Consumer groups divide partitions among members and a server-side offset records position, so a crash resumes exactly and a rewind replays.
No committed write lost when a broker dies
Replicate each partition to a leader plus in-sync replicas and commit only after acks=all, so a promoted ISR already holds every committed write.
Fail over in seconds without split-brain
A Controller runs leader election over a consensus store (ZooKeeper/KRaft) and propagates metadata so clients transparently reconnect.
A delivery guarantee apps can build on
At-least-once by default (retry until acked); idempotent producers + transactions layer exactly-once on top.

The trade-offs you say out loud

Senior signal isn’t the boxes — it’s naming what you gave up and why it was the right price.

An append-only commit log over a delete-on-read queue

Delete-on-read forgets too soon — no replay, and a second consumer can’t read what the first consumed; an append-only log keeps history and lets any number of consumers read independently.

Partitions routed by key hash over a bigger single broker

Vertical scaling hits one disk’s write rate and one CPU’s read rate; splitting a topic into partitions on many brokers scales throughput while a key-hash keeps related events ordered.

Consumer groups + server-side offsets over the broker pushing and deleting each message

Push-and-delete is back to a queue only one consumer can read with no replay; a stored offset cursor lets groups read the same log independently and rewind to replay.

Leader + in-sync replicas (acks=all) over periodic backups to cold storage

Backups are minutes stale, so a broker death still loses recent writes; committing only after in-sync replicas have the write means losing the leader loses nothing.

Controller-run leader election over clients electing a leader among themselves

Clients have no consistent global view and could split-brain; a controller promotes an in-sync replica via consensus — consensus is needed only for who-leads, not every message.

The answer, out loud

What a strong answer to “Design a Message Queue (like Kafka)” sounds like, first question to last trade-off. It is about 7 minutes of talking; the whiteboard and the interviewer fill the rest of the 45. Read it aloud once, then close the page and give it yourself.

  1. 0–3 min

    Scope it before drawing

    Let me scope it before I draw anything. The job is to sit between services that produce events and services that react to them, and the reason to put anything in the middle at all is that the two sides run at different speeds and fail at different times. So I want a producer to hand off and move on; I want several unrelated teams reading the same events without knowing about each other; and I want someone who joins next month, or who crashes today, to be able to go back and read history. Those three wants are what push me away from a queue and toward a log.

  2. 3–8 min

    The thing in the middle is a log

    The core is an append-only log. Writes go to the end and get an offset, and a read never removes anything. I like that because it makes the read side boring: a consumer is a number, not a subscription the broker has to remember deliveries for. It’s also the fast choice, since appending to the end of a file is sequential. The sentence I’d want the interviewer to hear is that this isn’t a transport, it’s storage — I’m choosing to keep the data and hand out cursors over it.

    Built in step 1: An append-only log
  3. 8–14 min

    Split it, and pick the key

    One log is one machine, so the next move is to split the topic into partitions, each its own log, each on a broker, with a partitioner choosing one by hashing a key. What I want to say plainly is what I gave up: order holds inside a partition, not across the topic. That’s usually fine, because what people actually need is order per entity — per user, per order — and the key buys exactly that. Choosing the key is the real design decision in this step, and I’d choose the entity whose events must never be reordered.

    Built in step 2: Partitions for parallelism
  4. 14–21 min

    How consumers read and share the work

    Then the read side. Instances of one service join a group, the group’s partitions are divided among the members, and a partition is handled by one member at a time, so work splits without anyone coordinating per message. Progress is an offset per group, stored on the broker side, so a restart resumes from a number instead of a guess. Two things fall out for free and I’d point at both: another team reads the same topic at its own position just by using its own group, and replay is only moving a number backwards.

    Built in step 3: Consumer groups & offsets
  5. 21–28 min

    Making a write actually durable

    Durability next. A partition on one disk is one disk away from gone, so each partition gets replicas — ×3 is the usual shape — with a leader taking reads and writes and the others following it. What makes this a guarantee rather than a hope is when I acknowledge, plus one thing I set next to it: with acks=all the producer isn’t told the write succeeded until the in-sync replicas have it, and I pair that with a floor on how small the in-sync set is allowed to get, because otherwise the word all can quietly come to mean the leader on its own. With that floor in place, a committed write really is one that more than one machine holds rather than one the leader acked by itself. I’m paying some latency per produce for that, and I’d take it on anything I’d be embarrassed to lose.

    Built in step 4: Replication
  6. 28–34 min

    Who decides who leads

    Now the leader dies. Someone has to notice and promote a replica, and I deliberately don’t let clients do it, because clients can’t agree and you end up with two brokers each certain it leads. Instead a controller watches broker health, elects a new leader from the in-sync set over a consensus store, and publishes the new metadata so producers and consumers reconnect on their own. The point I’d make is about where consensus goes: on who leads a partition, not on every message. That decision is small and rare, so the data path stays cheap.

    Built in step 5: Leader election
  7. 34–39 min

    Deciding what the log forgets

    The log can’t keep everything, so retention is a per-topic policy. For an event stream it’s age or size — keep 7 days, drop the older segments. For a changelog, where what I care about is the current value for each key, I’d compact instead and keep the latest record per key, which turns the log into a snapshot something can rebuild state from. I’d treat this as a product decision more than an operational one, because how far back the log goes is exactly how far back anyone can replay.

    Built in step 6: Retention & compaction
  8. 39–43 min

    What I would actually promise

    Last, the guarantee. The default is at-least-once: the producer retries until it gets an ack, and the consumer commits its offset after doing the work, so a crash re-delivers rather than drops. Duplicates are the price. Inside Kafka you can remove them — an idempotent producer dedupes retries, and transactions let the write and the offset commit land together — but the ordering of work and commit is the whole guarantee, and I’d want anyone on the team to be able to say which way round theirs is.

    Built in step 7: Delivery guarantees
  9. 43–45 min

    Close on the trade-off

    So: a log, partitioned by key, read by groups tracking offsets, replicated with acks=all, failed over by a controller, bounded by retention, at-least-once by default. The trade-off I’d name is that nearly everything good here comes from pushing state out to the reader — the broker stays dumb and fast, and the cost is that consumers own their own correctness. With more time I’d work through the rebalance in detail, and I’d size partition count against the peak consumer parallelism we really need, because that’s the number that hurts to change later.

What this teaches

Learn system design by building a distributed message queue like Kafka step by step. An interactive guide covering the append-only log, partitions, consumer groups and offsets, replication, leader election, retention, and delivery guarantees.

Key takeaways

  • An append-only commit log — ordered, immutable, replayable.
  • Partitions split the log across brokers; a key-hash keeps order per partition.
  • Consumer groups + server-side offsets: read once per group, rewind to replay.
  • Replication with a leader + in-sync replicas so no committed write is lost.
  • A controller runs leader election over a consensus store on broker failure.
  • Retention and compaction keep the log from growing without bound.
  • At-least-once by default; idempotent producers + transactions for exactly-once.

Concepts covered

  • What is a message queue?
  • An append-only log
  • Partitions for parallelism
  • Consumer groups & offsets
  • Replication
  • Leader election
  • Retention & compaction
  • Delivery guarantees
built to be replayed, not memorized — make the calls, kill the leader, run the gauntlet.
Finished this one? 0 / 65 System Designs done

Explore the topic

See this alongside everything else on the same subject — handbooks, system designs, challenges and tools, in one place.

More System Designs