1. What It Is
Kafka is the answer when multiple consumers need to replay the same event stream β not when you need a simple task queue. We think in topics, partitions, offsets, and consumer groups.
What:
A distributed, partitioned, replicated append-only event log. Brokers store topics split into ordered partitions; replication keeps data available when nodes fail.
Primary purpose:
High-throughput, durable asynchronous event streaming and ingestion.
Usually used for:
Event-driven pipelines, real-time analytics, decoupled microservices, and log aggregation.
2. Core Mental Model
Kafka is not a task queue β it is a replayable log. Think "who needs to read this stream, in what order, and for how long?" before drawing a broker box. Contrast with concept #17 Message Queues (consume-once tasks) and problem #32 Distributed Message Broker (full broker design interview).
β‘ Durable Event Pipeline
Unlike transient queues that delete messages after ACK, Kafka retains records on disk for a configurable window. Consumers read from an offset and can rewind.
π₯ Pull-Based Consumption
Consumers poll brokers on their own schedule. Slow consumers simply fall behind (consumer lag) instead of overwhelming the broker with push traffic.
π Partitioned Ordered Log
Concurrency is scaled by breaking topics into partitions. Ordering is strictly guaranteed only within a partition.
π‘οΈ Shock-Absorbing Async Layer
Sits between fast producers and slower downstream writers, absorbing bursts so spikes do not propagate as synchronous load.
In the room
The classic mistake is using Kafka like SQS. Emphasize ordering is per-partition only, consumers track offsets, and retention means replay is possible. If they ask about exactly-once, mention idempotent producers plus transactional writes β it's hard.
3. Why It Matters in HLD
Kafka fits when multiple consumers need the same durable, replayable stream β not for simple task queues. Three lenses:
Needed When:
You require high-volume data ingestion, async event decoupling, or replayable streams.
Avoids:
Tight synchronous coupling, cascading outages when one service slows down, and downstream databases overwhelmed by write spikes.
Optimizes For:
Write scalability, event ordering per entity, high-throughput durability, and reliable retries.
4. Architecture & Data Flow
Walk write and read paths as interview steps. Step 1 β Produce: producer batches records, picks partition via key hash. Step 2 β Replicate: leader appends to log; ISR followers sync before acks=all. Step 3 β Consume: consumer group polls batches, tracks offset per partition. Step 4 β Commit offset: after processing, commit β at-least-once by default. Step 5 β Scale: add partitions and consumers up to partition count; state ordering is per-partition only.
The Partition Model
Each topic consists of one or more partitions distributed across the cluster. Partitions allow Kafka to scale horizontally:
5. Key Characteristics
Partition keys, ISR, and pull-based consumption are the characteristics we cite on the whiteboard:
- Message record shape: Each record has a value (payload), optional key (partition routing), timestamp, and headers (metadata). Ordering is by offset within a partition, not by timestamp.
- Partition-Level Ordering: Messages are ordered strictly within a single partition, not globally across the topic.
- Partition key routing: Messages with the same key land on the same partition, preserving order per entity (e.g., per match, per user, per order).
- High Throughput via Sequential Writes: Kafka appends to partition logs sequentially, which maps well to disk and OS page cache.
- Zero-Copy Message Transfer: Data is copied directly from disk to network sockets without crossing into user-space memory.
- Offset-based progress: Consumers track their read position per partition and commit offsets back to Kafka. On restart they resume from the last commit (at-least-once by default).
- Decoupled Pull Protocol: Consumers request batches when they are ready, preventing memory exhaustion.
In the room
Emphasize ordering is per-partition only. If they ask exactly-once, mention idempotent producers plus transactional writes β then say it is hard and most teams aim for at-least-once with idempotent consumers.
6. Strategic Tradeoffs
Throughput and replay cost operational complexity and eventual consumer lag β we compare:
| Benefit | Cost |
|---|---|
| High Throughput (millions of events/sec via sequential I/O and zero-copy) | Operational Complexity (managing partitions, ZooKeeper/KRaft metadata, and client configs) |
| Durable Retries & Message Replay (allows historic data replay and easy recovery) | Eventual Consistency (readers might experience slight lag or delay across partitions) |
| Asynchronous Decoupling (isolates producers from consumer downtime/spikes) | Harder Debugging (tracing requests across async boundaries is complex) |
7. Failure / Bottleneck Awareness
Hot partitions, consumer lag, and rebalance pauses β we lead with mitigations before diagrams:
Problem: When key-based partitioning (e.g. partition by country code) sends a massive percentage of traffic to a single partition (like key: 'US'), overloading that broker while others sit idle.
Mitigation: Introduce composite partition keys (e.g. user_id + '_' + event_type) or add random salt suffixes to distribute heavy workloads evenly.
Problem: Downstream consumers process events slower than producers are publishing them, causing the consumer to fall further behind (lag grows). This can exhaust storage or delay business workflows.
Mitigation: Add more partitions to the topic and increase consumer instances within the Consumer Group (up to the partition limit) or optimize consumer processing with batching/multi-threading.
Problem: When a consumer joins or leaves the group (or crashes/fails to heartbeat), Kafka revokes all partition assignments, pauses processing, then reassigns partitions β inducing periodic latency spikes.
Mitigation: Tune heartbeat intervals (session.timeout.ms) and use Static Group Membership (Kafka 2.3+) to prevent rebalances during short networking blips or rolling deployments.
8. Common HLD Usage
These HLD problems map to Kafka when replay, fan-out, or per-entity ordering matter:
| Problem | Usage |
|---|---|
| Order Fulfillment Pipeline | Durable event log decoupling checkout from inventory, shipping, and notification workers |
| Real-time Activity Feeds (e.g. LinkedIn, Twitter) | Fanout event pipelines to processing workers |
| Uber Proximity / Dispatch matching | Geo-partitioned streaming pipeline to match riders and drivers |
| Application Metrics & Clickstream (Netflix, YouTube) | High-volume telemetry ingestion and log aggregation |
| Database Audit Logging / Change Data Capture (CDC) | Compacted state capture events streamed to search engines (Elasticsearch) |
9. Decision Signals
Choose Kafka over SQS/Rabbit when replay, multiple independent consumers, or high-volume ingest dominate:
- You need async, fire-and-forget workflows with durable delivery (tune producer
acksand replication for your SLA). - You must replay and re-process historic messages (e.g. rebuilding read-model caches).
- You require event ordering strictly per entity (e.g. tracking banking transaction ledger events).
- You face high-volume streaming ingest (clickstreams, metrics, IoT sensors).
- Multiple independent microservices need to consume the exact same stream of events.
11. Deep Dive (Optional)
ISR & Durability Mathematics
Kafka achieves partition reliability via replication. The leader replica handles producer writes and consumer reads (by default). Followers that are caught up with the leader form the ISR (In-Sync Replicas) set; lagging replicas drop out of acks=all quorum:
To balance write speed against durability risk, tune the producer's acks configuration:
acks=0: Producer fire-and-forget. Zero durability guarantee.acks=1: Returns success the instant the leader broker writes it locally. Risk: leader dies before replication.acks=all: Leader writes locally and waits for all active ISR followers to sync. Full durability, higher write latency.
Log Compaction Internals
Normally, Kafka purges logs by age (TTL) or log size limit. But for key-value streams (like database records captured via CDC), we only care about the latest state of a key. Log Compaction periodically garbage-collects old keys, keeping only the final updated version per key:
Zero-Copy Reads & OS Page Cache
Traditional brokers read file contents into kernel buffer, copy them to user-space application memory, copy them back to kernel socket buffer, and then to network card. Kafka bypasses this entirely using the sendfile system call. The OS transfers data directly from the kernel page cache to the network NIC, eliminating CPU overhead and context switches.
ZooKeeper vs KRaft Consensus
Historically, Kafka relied on ZooKeeper to manage cluster states, partition assignments, and broker membership. Under KRaft (modern Kafka), consensus is integrated directly into Kafka brokers using a customized Raft-based metadata log. This eliminates the external coordination dependency, allows instant cluster boots, and supports millions of partitions.
Retention vs Compaction: Know Your Stream Type
Kafka offers two log lifecycle modes β choosing wrong breaks downstream consumers:
- Time/size retention (delete): Old segments drop after N days or N GB. Correct for clickstreams, metrics, and audit logs where history beyond the window is worthless.
- Log compaction: Keeps the latest record per key forever; tombstones with null values delete keys. Correct for CDC changelog topics where consumers rebuild current state from the log.
Event sourcing pitfall: Never enable compaction on an append-only event-sourcing topic where every state transition must be preserved β compaction garbage-collects intermediate events and destroys the audit trail. Pair event sourcing with infinite retention (delete policy disabled) or external snapshot + archive, not key compaction.
Review
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.