๐จ 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_metadatalog, 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 whenbatch.sizeis reached orlinger.msexpires- 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_committedconsumers 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 withacks=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