Skip to main content

Kafka Broker — Complete Guide

Kafka Broker — Storage Engine InternalsClick a component to inspect
KAFKA BROKER NODE — Log Storage Layer/data/orders-0/00000000.log00000000.index00000000.timeindex00001048.log (active)/data/orders-1/00000000.log00000000.index00000000.timeindex00001048.log (active)/data/orders-2/00000000.log00000000.index00000000.timeindex00001048.log (active).logRecordBatches.indexOffset → PosSender ThreadNIO non-blocking I/OReplicaFetcher ThreadFetchRequest → followerproduce writereplica fetch
Click a component to inspect internals

Kafka Binary RecordBatch Structure & Segment Index Layout

Kafka Binary RecordBatch (v2) & Memory Offset Index Architecture
Click a binary field to inspect its payload encoding and role:
BaseOffsetLogical Addressing
Memory Overhead
8 Bytes
Wire Data Type
Int64
Starting logical offset assigned by broker to the first record in this batch.

Who this guide is for

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

ResponsibilityArchitectural Mechanism
Message StorageAppends incoming record batches to append-only .log segment files on disk.
Partition LeadershipHandles 100% of producer writes and consumer reads for its assigned leader partitions.
Data ReplicationFetches records from leader partitions as a follower to maintain configured Replication Factors.
Offset ManagementManages consumer group commit offsets inside the internal __consumer_offsets system topic.
Metadata CoordinationCommunicates 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
  1. Segment Log (.log): Stores serialized Kafka RecordBatch structures containing payloads, headers, timestamps, and sequence numbers.
  2. Offset Index (.index): A sparse index mapping logical record offsets to exact physical byte positions inside the .log file. Instead of indexing every record, Kafka inserts an entry every 4 KB4\text{ KB} (index.interval.bytes).
  3. 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 ≈200,000\approx 200,000 due to znode memory and watch notification overhead.
  • Failover Latency: When a Controller broker crashed, electing a new Controller and repopulating metadata took 15–30 seconds15\text{--}30\text{ seconds}.
  • KRaft Advantage: Stores metadata directly in a specialized internal partition (__cluster_metadata). KRaft Controller failovers execute in sub-second timeframes, supporting clusters with over 1,000,0001,000,000 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 (300–600 MB/sec300\text{--}600\text{ MB/sec} on SATA SSDs, 3–7 GB/sec3\text{--}7\text{ GB/sec} 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

Kafka Zero-Copy: sendfile(), Page Cache & DMA Engine Deep Dive
USER SPACE (Ring 3: JVM Process / Kafka Broker)KERNEL SPACE (Ring 0: Linux OS Page Cache & Sockets)HARDWARE CONTROLLERS (DMA & NIC)NVMe / DiskPartition .logOS Page CacheKernel RAMSocket DescriptorOffset & Length pointer only⚡ Direct Scatter-Gather DMA Highway (0 CPU Memory Copies!)NIC BufferConsumer Network
CPU Memory Copies
0 Copies
Context Switches
2 Switches
CPU Overhead
Near Zero (Pure DMA)
JVM Heap Footprint
0 MB (Bypasses JVM Heap)
LIFECYCLE STEPS: Linux sendfile() / Java FileChannel.transferTo()
1. sendfile() Syscall
Kafka broker executes FileChannel.transferTo(). CPU switches from User Mode to Kernel Mode (Context Switch 1).
2. DMA Read to Page Cache
DMA copies data from Disk into OS Page Cache (or reads existing cached RAM directly). No CPU copy needed.
3. File Descriptor Append
Only small descriptor metadata (pointer + length) is passed to the socket buffer. No payload data is copied!
4. Scatter-Gather DMA to NIC
NIC hardware directly gathers payload bytes from OS Page Cache memory via DMA. CPU switches back to User Mode (Context Switch 2).
SENIOR ARCHITECTURE INSIGHT

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.

JVM Heap Sizing Rule: Keep broker heap small (-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 →\to OS Page Cache (DMA) →\to JVM Heap Buffer (CPU copy + context switch) →\to Kernel Socket Buffer (CPU copy + context switch) →\to 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 pure sendfile() 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

MetricTarget ValueAlert ConditionOperational Meaning
UnderReplicatedPartitions0>0> 0Partitions have lost ISR members — data loss risk.
ActiveControllerCount1≠1\neq 1Cluster has no controller or is in split-brain state.
OfflinePartitionsCount0>0> 0Partitions have no active leader — reads/writes failing.
BytesInPerSecNominal baselineApproaching NIC bandwidthNetwork ingress saturation.
RequestHandlerAvgIdlePercent>50%> 50\%<30%< 30\%Broker handler thread pool is exhausted.

Common Failure Scenarios

FailureRoot CauseRemediation
Broker OOM KillJVM heap sized too large, starving OS Page Cache and native memory.Cap JVM heap at 6 GB6\text{ GB} (-Xmx6g); reserve remaining host RAM for Linux Page Cache.
ISR FlappingReplica 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 LossOut-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_offsets topic.

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 with acks=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's FileChannel.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.


See Also

📖
Track Page Progress0 / 635 Read
Knowledge Base Completion0%