Kafka Broker — Complete Guide
Kafka Binary RecordBatch Structure & Segment Index Layout
- New learners — start at What is Kafka? and What is a Broker? to understand the foundational model before diving into storage and replication.
- Senior engineers — jump to Replication & ISR, Broker Internals, KRaft vs ZooKeeper, Log Compaction, or Performance Tuning.
What is a Kafka Broker?
A Kafka Broker is a dedicated server process executing within an Apache Kafka cluster. The broker's primary duty is to receive incoming event streams from producers, append them to sequential partition log segment files on disk, serve messages to consumer groups, manage data replication across cluster nodes, and maintain cluster metadata state.
Broker Responsibilities
| Responsibility | Architectural Mechanism |
|---|---|
| Message Storage | Appends incoming record batches to append-only .log segment files on disk. |
| Partition Leadership | Handles 100% of producer writes and consumer reads for its assigned leader partitions. |
| Data Replication | Fetches records from leader partitions as a follower to maintain configured Replication Factors. |
| Offset Management | Manages consumer group commit offsets inside the internal __consumer_offsets system topic. |
| Metadata Coordination | Communicates with the cluster Controller node (via KRaft or ZooKeeper) to manage topology changes. |
Core Broker Storage Engine Internals
Kafka stores partition records on disk in dedicated topic-partition directories (/var/lib/kafka/data/<topic>-<partition_id>/). Each partition directory contains a sequence of Log Segments:
/var/lib/kafka/data/orders-0/
├── 00000000000000000000.log # Raw binary RecordBatches
├── 00000000000000000000.index # Sparse offset -> physical position index
├── 00000000000000000000.timeindex # Timestamp -> offset index
├── 00000000000000001048.log # Active segment currently accepting writes
├── 00000000000000001048.index
└── leader-epoch-checkpoint # Tracks leader epoch history for fencing
- Segment Log (
.log): Stores serialized KafkaRecordBatchstructures containing payloads, headers, timestamps, and sequence numbers. - Offset Index (
.index): A sparse index mapping logical record offsets to exact physical byte positions inside the.logfile. Instead of indexing every record, Kafka inserts an entry every (index.interval.bytes). - Time Index (
.timeindex): Maps timestamps to logical offsets, supporting time-based consumer seeking (KafkaConsumer.offsetsForTimes()).
KRaft (Kafka Raft) Metadata Architecture
Historically, Kafka relied on external Apache ZooKeeper clusters for metadata storage and active controller elections.
Legacy ZooKeeper Mode: Modern KRaft Mode (Kafka 3.3+):
+--------------------+ +------------------------------------+
| ZooKeeper Quorum | | Active KRaft Controller Broker |
| (External 3 nodes) | | - Stores __cluster_metadata Log |
+---------+----------+ +-----------------+------------------+
| |
v v Quorum Replication
+--------------------+ +-----------------+------------------+
| Kafka Controller | | Broker 2 | Broker 3 |
| (Broker Node 1) | | (Controller) | (Controller) |
+--------------------+ +-----------------+------------------+
Why ZooKeeper Was Removed
- Scale Limit: ZooKeeper capped total cluster partition counts at due to znode memory and watch notification overhead.
- Failover Latency: When a Controller broker crashed, electing a new Controller and repopulating metadata took .
- KRaft Advantage: Stores metadata directly in a specialized internal partition (
__cluster_metadata). KRaft Controller failovers execute in sub-second timeframes, supporting clusters with over partitions.
Performance Mechanics: Why Kafka is Fast
1. Sequential Disk I/O vs Random Access
Kafka performs append-only sequential writes to log segments. While random disk access on spinning disks or NVMe drives incurs heavy seek latency and IOPS limits, sequential disk throughput ( on SATA SSDs, on NVMe) approaches raw hardware bus throughput.
2. OS Page Cache vs In-JVM Caching
Kafka does not cache message data in the JVM heap. Storing objects in JVM memory causes double caching (once in the OS Page Cache and once in the JVM), high GC pause overhead, and 2-4x memory bloat due to Java object headers. Instead, all disk reads and writes flow directly through the Linux OS Page Cache in kernel RAM. Recent records are read directly from RAM by consumers at memory bus speeds without touching physical disk platters or flash cells.
3. Zero-Copy sendfile() Syscall & Scatter-Gather DMA
Why Page Cache + sendfile() Beats In-Process Caching
By delegating caching entirely to the Linux OS Page Cache, Kafka avoids object serialization overhead, GC pauses, and memory fragmentation inside the JVM. When multiple consumers read the same partition (e.g. 10 consumers in different groups), Linux serves all 10 consumers directly from Page Cache RAM via DMA with 0 CPU memory copies.
-Xms6g -Xmx6g). Leave 75%+ of physical host RAM free for the Linux Page Cache to saturate 100GbE NICs effortlessly.When a consumer sends a FetchRequest, the broker executes Java's FileChannel.transferTo(), which maps directly to the Linux sendfile() system call:
- Traditional I/O (4 Copies, 4 Context Switches): Data moves from Disk OS Page Cache (DMA) JVM Heap Buffer (CPU copy + context switch) Kernel Socket Buffer (CPU copy + context switch) NIC Buffer (DMA). Copying 10 Gbps of traffic requires significant CPU saturation solely moving memory bytes.
- Linux Zero-Copy
sendfile()(0 CPU Copies, 2 Context Switches): The broker passes only socket descriptor metadata (offset + length) to the kernel socket buffer. The network interface card (NIC) uses Scatter-Gather DMA to gather payload bytes directly from the OS Page Cache RAM and stream them across the network cable. - TLS / SSL Impact & Kernel TLS (kTLS): Standard user-space TLS encryption (Java
SSLEngine) breaks puresendfile()zero-copy because data must pass through the JVM to be encrypted. High-throughput clusters leverage Kernel TLS (kTLS) (Linux 4.17+) or TLS hardware offload to restore zero-copy throughput under encryption.
Monitoring & Observability Matrix
| Metric | Target Value | Alert Condition | Operational Meaning |
|---|---|---|---|
UnderReplicatedPartitions | 0 | Partitions have lost ISR members — data loss risk. | |
ActiveControllerCount | 1 | Cluster has no controller or is in split-brain state. | |
OfflinePartitionsCount | 0 | Partitions have no active leader — reads/writes failing. | |
BytesInPerSec | Nominal baseline | Approaching NIC bandwidth | Network ingress saturation. |
RequestHandlerAvgIdlePercent | Broker handler thread pool is exhausted. |
Common Failure Scenarios
| Failure | Root Cause | Remediation |
|---|---|---|
| Broker OOM Kill | JVM heap sized too large, starving OS Page Cache and native memory. | Cap JVM heap at (-Xmx6g); reserve remaining host RAM for Linux Page Cache. |
| ISR Flapping | Replica network lag or G1GC pauses exceeding replica.lag.time.max.ms. | Tune G1GC pause targets (-XX:MaxGCPauseMillis=20) and check disk latency (iostat -xz 1). |
| Unclean Leader Election Data Loss | Out-of-sync follower elected as leader when all ISR members crash. | Set unclean.leader.election.enable=false in production. |
Interview Questions
Q1. What is a Kafka broker and what are its primary responsibilities?
A Kafka broker is a single server node running the Kafka process. Its duties include: appending records to sequential partition log segment files on disk, serving fetch requests from consumers, executing follower replication from leader brokers, and storing consumer group offsets in the
__consumer_offsetstopic.
Q2. What is the In-Sync Replicas (ISR) set and why is it critical for consistency?
The ISR set consists of all partition replicas (leader + followers) currently caught up with the leader within
replica.lag.time.max.ms. When a producer writes withacks=all, the leader waits for acknowledgements from all current ISR members before returning success. If the leader fails, only an ISR member is eligible for leader election, guaranteeing zero data loss.
Q3. How does Zero-Copy data transfer work in Kafka?
Zero-Copy leverages the Linux
sendfile()system call (exposed via Java'sFileChannel.transferTo()). Data stored in the OS Page Cache is transferred directly to the Network Interface Card (NIC) buffer via Direct Memory Access (DMA), bypassing the JVM heap and eliminating user-kernel context switches.
Q4. What is the difference between log.retention.hours (deletion) and log.cleanup.policy=compact?
Log deletion (
delete) purges old log segments once they exceed time or size thresholds. Log compaction (compact) retains the single latest record payload for every message key indefinitely while deleting earlier historical updates for that key. Compaction is ideal for event-sourced state and lookup changelogs.
