SysDesignPrep.com
Study guide 87 of 183

How Kafka works

Kafka's internals for system design interviews: topics, partitions and the append-only log, segments and indexes, replication with leaders and in-sync replicas, acks and durability, consumer groups and offsets, retention and compaction, ordering and throughput.

Reading is half of it. See this used in a real interview: walk through Design a Distributed Message Queue →

Kafka is the default event streaming platform in system design answers, and "Design a Message Queue" is often really "explain how Kafka works". Its design is surprisingly simple: a partitioned, replicated, append-only log, read by consumers that track their own position. That simplicity is why it is so fast and why it supports replay, multiple independent consumers and very high throughput.

Topics, partitions and the log

  • A topic is a named stream (orders, clicks).
  • Each topic is split into partitions; each partition is an ordered, append-only log on disk.
  • Every record in a partition has an offset, its position in the log.
  • Producers choose the partition, usually by hashing a key (all events for order 42 go to the same partition, preserving their order).

Ordering is guaranteed within a partition only. Partitions are the unit of parallelism: more partitions allow more consumers and more throughput.

Storage

Each partition is a directory of segment files (for example 1 GB each) plus small indexes from offset (and timestamp) to file position. Writes append to the active segment, which is sequential I/O. Reads are served largely from the operating system's page cache, and Kafka uses zero-copy transfer to send data from disk to network without passing through application memory. Records are batched and compressed per batch. This is why a modest cluster handles millions of messages per second. See storage engines.

Replication

Each partition has a leader replica and followers on other brokers (replication factor typically 3):

  • Producers write to the leader; followers fetch from it.
  • The in-sync replica set (ISR) contains followers that are caught up.
  • A record is committed once all ISR members have it; consumers only see committed records.
  • If the leader fails, a new leader is elected from the ISR (coordinated by the controller, using KRaft consensus in modern Kafka). See consensus and coordination.

Durability settings

SettingEffect
acks=0do not wait; fastest, may lose data
acks=1leader has it; lost if the leader dies before followers copy it
acks=all with min.insync.replicas=2at least two replicas have it; survives one broker loss
Idempotent producerretries do not create duplicates

For important data use acks=all, min.insync.replicas=2, replication factor 3 and idempotent producers. Note that Kafka relies on replication more than fsync for durability.

Consumers and consumer groups

  • Consumers pull records and track their offset per partition.
  • A consumer group shares the work: each partition is assigned to exactly one consumer in the group, so a group can have at most as many active consumers as partitions.
  • Different groups read independently: the search indexer, the analytics pipeline and the notification service each read the full stream at their own pace.
  • Offsets are committed to an internal topic; committing after processing gives at-least-once semantics. See delivery semantics.
  • When consumers join or leave, partitions are rebalanced; cooperative rebalancing avoids stopping the whole group.

Because messages are not deleted when read, consumers can replay from any retained offset: reprocess after a bug fix or bootstrap a new service from history.

Retention and compaction

  • Time or size retention: delete segments older than 7 days (or any period). Storage is cheap, so long retention is common; tiered storage moves old segments to object storage.
  • Log compaction: keep only the latest record per key (and delete keys with tombstones). The topic becomes a durable, replayable table of current state, used for changelogs and configuration. See change data capture.

Sizing partitions

  • Throughput: if one partition handles around 10 MB/s for your consumers and you need 200 MB/s, you need at least 20 partitions, plus room to grow.
  • Parallelism: enough partitions for the maximum number of consumers you will run.
  • Not too many: each partition costs memory, file handles and failover time. Increasing partitions later changes key-to-partition mapping, which breaks ordering for existing keys, so choose with headroom.
  • Watch for hot partitions from skewed keys. See hot keys and skew.

Kafka versus a traditional queue

KafkaRabbitMQ or SQS-style queue
Modellog; consumers track offsetsmessages removed when acknowledged
Replayyes, within retentionno
Multiple consumersindependent groups read everythingcompeting consumers share messages
Orderingper partitionlimited (FIFO queues per group)
Per-message featuresbasicdelays, priorities, per-message routing and retries
Strengthhigh throughput event streams, fan-out, replaytask queues, per-message acknowledgement

See message queues and streams.

In the interview

For Design a Message Queue: partitioned append-only logs in segments, keyed partitioning for order, leader-follower replication with an ISR and acks=all, consumer groups with committed offsets, retention and compaction, and partition count chosen from throughput and parallelism. For pipelines like Design an Ad Click Aggregator, mention replay and exactly-once processing.

Checklist

  • Partitions as the unit of order, parallelism and throughput.
  • Sequential segment writes, page cache, batching and compression.
  • Replication factor 3, acks=all, min.insync.replicas=2, idempotent producers.
  • Consumer groups, offsets committed after processing, rebalancing.
  • Retention for replay; compaction for latest state per key.
  • Partition count sized with headroom; hot keys watched.

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.