On this page
Tracks

System Design — Distributed Message Queue

Last reviewed 11 Sept 2026

Part of the system design series. See the framework and building blocks first if you haven’t.

1. Requirements

Functional

  • Producers publish messages to named topics; consumers subscribe and read them, independently and at their own pace.
  • Multiple consumers must be able to process a topic in parallel (consumer groups) without duplicating work.
  • Messages within a partition are delivered in the order they were written.
  • Support replay — a consumer can rewind and reprocess messages, unlike a fire-and-forget queue.

Non-functional

  • High sustained write throughput (hundreds of thousands to millions of messages/sec at large-scale deployments).
  • Durability — an acknowledged write must survive a broker crash.
  • Horizontal scalability for both storage and throughput as topics grow.
  • Configurable delivery/ordering guarantees, because “exactly-once, globally ordered, infinitely durable, infinitely fast” cannot all be had at once — the interview is about which three you pick and why.

2. Where it sits / high-level architecture

flowchart TD
P1[Producer] -->|key hash -> partition| B1[Broker 1 - Partition 0 leader]
P1 --> B2[Broker 2 - Partition 1 leader]
B1 -->|replicate| B3[Broker 3 - Partition 0 follower]
B2 -->|replicate| B1b[Broker 1 - Partition 1 follower]
C1[Consumer group A - member 1] -->|reads Partition 0| B1
C2[Consumer group A - member 2] -->|reads Partition 1| B2
C3[Consumer group B - independent offset] -->|reads Partition 0| B1
CTRL[Controller / metadata quorum] -.leader election, membership.-> B1
CTRL -.-> B2
CTRL -.-> B3
Topic split into ordered partitions, each replicated across brokers
  • A topic is a logical stream split into partitions, each an append-only, ordered log stored on disk — the partition, not the topic, is the unit of parallelism, ordering, and storage.
  • Each partition has one leader broker taking all reads/writes for it, and a set of follower replicas that continuously fetch and replay the leader’s log for durability.
  • A controller (in modern Kafka, a small quorum of nodes running the Raft-based KRaft protocol, replacing the old ZooKeeper-based control plane) tracks broker membership, partition leadership, and triggers re-election on failure.
  • Consumer groups let multiple consumers split a topic’s partitions among themselves — each partition is read by exactly one consumer within a group at a time, which is what turns “ordered log” into “scalable parallel processing” without losing per-partition order.

3. Ordering, partitioning, and delivery-guarantee trade-offs

ChoiceWhat you getWhat you give up
More partitionsMore parallelism (more consumers can work concurrently), higher aggregate throughputWeaker ordering guarantee (only per-partition, never topic-wide); more replication/metadata overhead; slower leader election on failure since more partitions must fail over
Fewer partitionsSimpler, tighter ordering, less overheadCaps how many consumers can usefully parallelize (a consumer group can’t have more active consumers than partitions — extras sit idle)
Key-based partitioning (hash(order_id) % N)All events for the same key stay in order (e.g. all events for one order)Hot keys create hot partitions — a viral key can overload one partition while others idle
Round-robin / no keyEven load distributionNo ordering guarantee across events that should logically relate
At-most-once (fire and forget)Lowest latency, simplestSilent message loss on any failure
At-least-once (ack after persist, retry on doubt)No message lossConsumer must handle duplicates — needs idempotent processing
Exactly-once semantics (Kafka transactions / idempotent producer + transactional consumer)Strongest guarantee, no dupes, no lossReal throughput and latency cost; genuinely complex to reason about across producer retries and consumer offset commits

Partition count sizing — pick partitions based on your target parallelism and per-partition throughput ceiling (a single partition on typical hardware handles on the order of low tens of MB/s), not an arbitrary round number; over-partitioning a topic “for future scale” has a real cost in controller metadata and failover time.

4. Deep dive — replication (ISR) and consumer group rebalancing

In-sync replica (ISR) set. Each partition’s leader tracks which followers are “caught up” (within a configurable lag threshold) — this set is the ISR. A producer configured for durability (acks=all) only considers a write committed once every replica in the ISR has it, not just the leader. If a follower falls behind (slow disk, network blip), the controller shrinks the ISR to exclude it — this keeps the durability guarantee honest: “all in-sync replicas have it” only means something if replicas that fall behind get removed from the set rather than silently counted as safe. When the lagging follower catches back up, it re-joins the ISR. This dynamic shrink/grow is the mechanism that lets Kafka tolerate slow or flaky followers without either blocking all writes on the slowest node or silently weakening the durability guarantee.

Leader election on failure. When a partition’s leader broker dies, the controller detects it (via the KRaft quorum’s failure detection) and promotes an in-sync follower to leader. Only ISR members are eligible by default — this is the guarantee that keeps you from promoting a replica that’s missing recent writes and silently losing “committed” data. A partition with no leader is unavailable for reads/writes on that partition specifically until election completes; other partitions on other brokers are unaffected, which is why partition count vs. failover-time is a real trade-off (more partitions per broker means more individual elections to run when that broker dies).

Consumer group rebalancing. When a consumer joins, leaves, or is detected as dead (missed heartbeats), the group coordinator reassigns partitions among the remaining members. Naive “stop-the-world” rebalancing (every consumer pauses while reassignment happens) causes a processing pause proportional to group size; modern Kafka supports incremental cooperative rebalancing, where only the specific partitions that need to move are revoked and reassigned, letting unaffected consumers keep processing throughout.

5. What real systems do today

  • Apache Kafka’s KRaft mode (replacing ZooKeeper as of recent major versions) moves cluster metadata and controller election onto a Raft-based quorum built into Kafka itself, removing the operational overhead of running a separate ZooKeeper ensemble — this is the current recommended architecture as of 2025-2026 deployments.
  • Leader-based replication with an ISR set, as described above, is the mechanism Kafka and Confluent’s managed offering both document as the core durability guarantee — acks=all plus min.insync.replicas is the standard production-durability configuration pair engineering teams tune for “don’t lose committed data.”
  • Partition-count-as-parallelism-unit is repeatedly the top piece of advice in current Kafka topic-design guides: pick partition count from target throughput and consumer parallelism, and treat it as a decision that’s expensive to change later (you can increase partitions, but it breaks the ordering guarantee for existing keys since the hash-to-partition mapping shifts).
  • Incremental cooperative rebalancing is the modern default over the older “revoke everything, reassign everything” protocol specifically because full-group pauses on every membership change were a recurring operational pain point at high consumer-group churn (e.g. rolling deploys of a consumer fleet).
  • Monitoring guidance from observability vendors converges on watching ISR shrink events, under-replicated partitions, and consumer lag as the primary health signals for a Kafka cluster — these are the metrics that catch replication and consumption problems before they become customer-visible.

6. Scaling & failure

BottleneckFixNew cost
Single partition throughput ceiling reached (hot key or genuinely high volume)Increase partition count, or split a hot key into sub-keys (order_id + shard_suffix)Existing consumers must rebalance; ordering guarantee is now per-sub-key, not per original key
Broker running out of disk from long retentionTiered storage — move older log segments to cheap object storage, keep only recent hot segments on local diskSlightly higher read latency for old data, but retention becomes effectively unbounded and cheap
Consumer group falling behind (lag growing)Add more consumers up to the partition count, or speed up per-message processing (batch, async I/O)Consumers beyond partition count sit idle — this is the hard ceiling, not compute
Controller quorum metadata operations slow under very high partition countsKeep partition count proportional to actual need, monitor controller metadata sizeRequires periodic partition-count review as topics evolve, not a “set once” decision

What happens when the controller quorum loses majority

  • With KRaft, the controller role itself is a Raft-replicated quorum (typically 3 or 5 nodes). If enough controller nodes are down that the quorum can’t reach majority (e.g. 2 of 3 down), no new controller can be elected and no cluster metadata changes can happen — no partition leader reassignment, no topic creation, no ISR updates.
  • Critically, partitions whose current leader is still alive and reachable keep serving reads and writes normally — the data plane and control plane are decoupled. The failure mode is that the cluster becomes “frozen” for anything requiring a metadata change: if a data-plane broker also dies during this window, its partitions become stuck without a leader until the controller quorum recovers, because nobody can run the election.
  • This is the sharp edge to name explicitly in an interview: a Kafka cluster degrades gracefully on a single broker or a single controller-quorum-member loss, but a simultaneous loss of controller quorum majority plus a data-plane broker is the compounding failure that actually causes an outage — which is exactly why controller nodes are provisioned with extra redundancy and spread across failure domains separately from data brokers.

Interview follow-ups

  • “Why is ‘exactly-once delivery’ the wrong thing to promise, and what do you promise instead?” — Network ambiguity between “ack lost” and “message lost” makes true exactly-once delivery impossible; promise at-least-once delivery plus idempotent consumer-side processing (dedupe by message id).
  • “How does adding partitions affect ordering for existing data?” — The key-to-partition hash changes, so messages for the same key can land on a different partition going forward — ordering is only ever guaranteed within a partition, and repartitioning breaks continuity for in-flight keys.
  • “Walk me through what happens when a partition’s leader broker dies.” — Controller detects it via the KRaft quorum, promotes an ISR member (never a replica that’s behind) to leader; that partition is briefly unavailable during election, other partitions unaffected.
  • “What’s the difference between a replica and an in-sync replica?” — ISR membership requires being caught up within a lag threshold; a lagging replica is dropped from ISR so acks=all durability claims stay honest, and rejoins once it catches up.
  • “A consumer in a group crashes. What happens to the partitions it was reading?” — Heartbeat timeout triggers rebalancing; with cooperative incremental rebalancing, only the orphaned partitions move, other consumers in the group keep processing uninterrupted.
  • “How would you handle a single hot key overwhelming one partition?” — Sub-key sharding (append a shard suffix) trades strict per-original-key ordering for load distribution — name the trade-off explicitly.
  • “What’s the actual difference between this and a traditional queue like SQS/RabbitMQ?” — A queue typically removes a message once consumed (single logical consumer per message); a log-based system like Kafka retains messages for a retention window so multiple independent consumer groups can each read the full stream at their own offset, enabling replay.

Sources: Apache Kafka architecture: A complete guide 2026 — Instaclustr · KRaft Overview — Confluent Documentation · How to Design Kafka Topics and Partitions — OneUptime · Kafka Deep Dive for System Design Interviews — Hello Interview · Kafka: Topics, Partitions, Replication, ISR, Leader Election, Acks — Deep Dive · Monitoring Kafka performance metrics — Datadog