Core Concept

Kafka Architecture and Guarantees

Kafka is a distributed append-only log for high-throughput event streaming. Producers write to partitioned topics; consumer groups pull at their own pace and can replay history when retention allows.


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.

Loading...

The Partition Model

Each topic consists of one or more partitions distributed across the cluster. Partitions allow Kafka to scale horizontally:

Loading...

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:

BenefitCost
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:

πŸ”₯ Hot Partitions (Skew)

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.

Loading...
🐒 Consumer Lag

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.

Loading...
πŸ›‘ Rebalance Pauses (Stop-The-World)

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.

Loading...

8. Common HLD Usage

These HLD problems map to Kafka when replay, fan-out, or per-entity ordering matter:

ProblemUsage
Order Fulfillment PipelineDurable 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 matchingGeo-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:

🎯 Think Kafka When:
  • You need async, fire-and-forget workflows with durable delivery (tune producer acks and 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:

Loading...

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:

Loading...

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

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...