๐Ÿ“จ Kafka Internals (important parts only)

1. Storage: the log

  • Topic โ†’ partitions โ†’ each partition is a directory of segments: .log (records in batches), .index (sparse offset โ†’ file position), .timeindex
  • Append-only, sequential writes; reads mostly hit the OS page cache; zero-copy (sendfile) to consumers when not using TLS
  • Retention by time/size deletes whole old segments; log compaction keeps the latest value per key (a cleaner thread; tombstones = null values)
  • Why Kafka is fast: sequential I/O + page cache + batching + compression per batch + zero-copy + partition parallelism

2. Replication

  • Each partition has a leader and followers; followers pull (fetch) from the leader
  • ISR (in-sync replicas): followers caught up within replica.lag.time.max.ms
  • LEO (log end offset) per replica; high watermark (HW) = the highest offset replicated to all ISR โ†’ consumers only see up to the HW
  • acks=all + min.insync.replicas=2 (RF=3) โ†’ a write succeeds only when โ‰ฅ 2 replicas have it
  • Leader epochs make followers truncate divergent logs after a leader change
  • unclean.leader.election.enable=false โ†’ prefer unavailability over data loss

3. Control plane: KRaft

  • Metadata (topics, partitions, ISR, configs) lives in the __cluster_metadata log, replicated by Raft across controller nodes; brokers replicate that log. ZooKeeper is gone as of Kafka 4.0

4. Producer

  • send() โ†’ serializer โ†’ partitioner (murmur2(key) % partitions; sticky batching for null keys) โ†’ RecordAccumulator (a batch per partition) โ†’ sender thread ships batches when batch.size is reached or linger.ms expires
  • Idempotent producer (default on): producer ID + a per-partition sequence number โ†’ the broker drops duplicates from retries; ordering is preserved with โ‰ค 5 in-flight requests
  • Transactions: transactional.id โ†’ a transaction coordinator (__transaction_state); writes to many partitions + consumer offsets commit atomically; control records (commit/abort markers); read_committed consumers read up to the LSO (last stable offset)

5. Consumer

  • A pull-based poll() loop; fetch.min.bytes/fetch.max.wait.ms; max.poll.records
  • Liveness: a heartbeat thread (session.timeout.ms) and poll progress (max.poll.interval.ms); slow processing โ†’ kicked from the group โ†’ rebalance
  • Group coordinator; offsets stored in __consumer_offsets (compacted)
  • Rebalancing: eager (stop-the-world) โ†’ cooperative incremental โ†’ the new consumer group protocol (KIP-848, server-side assignment)
  • Delivery: commit after processing = at-least-once โ†’ make handlers idempotent

๐Ÿ”ฌ Prove it

  • 3-broker KRaft cluster in Docker; kafka-topics --describe โ†’ kill the leader โ†’ watch the ISR shrink and a new leader get elected
  • Data loss experiment: acks=1, kill the leader right after writes โ†’ count lost messages; repeat with acks=all, min.insync.replicas=2 โ†’ zero loss
  • Inspect segments with kafka-dump-log.sh --files โ€ฆ --print-data-log; find batches, compression, and producer IDs
  • Throughput vs linger.ms (0, 5, 20) and compression (none/lz4/zstd): plot it
  • Slow consumer exceeding max.poll.interval.ms โ†’ observe the rebalance โ†’ fix it with smaller batches or pause/resume
  • Ordering test: key = runId across 6 partitions, scale consumers 1 โ†’ 3 โ†’ 6
  • Transactions: read-process-write with read_committed; abort midway and show consumers never see the aborted records

Interview questions interview-q

How Kafka guarantees ordering ยท the HW vs LEO ยท acks + min.insync.replicas trade-offs ยท how idempotent producers work ยท exactly-once in Kafka (and its limits outside Kafka) ยท why Kafka is fast ยท rebalancing problems ยท Kafka vs RabbitMQ