skip to content

Architecture and Internals

How a Kafka broker actually stores and serves data: append-only segment files, indexes, page cache, and the request-handling threads. Interviewers dig here to separate people who have run Kafka from people who have only used the client API.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

explore

questions

page 2 of 2

Walk through how a closed log segment gets offloaded to remote storage. What does the broker upload, and when can it delete the local copy?

level: middleimportance: should knowfreq 50%

basics

~20 s

When a segment rolls (closes), the leader broker's RemoteLogManager uploads it plus its indexes to remote storage and records metadata. Only after the upload is confirmed (COPY_SEGMENT_FINISHED) can the local copy be deleted, subject to local retention.

open as a page

How does broker.rack enable rack awareness, and what does it (and doesn't it) guarantee for replica placement and reads?

level: seniorimportance: should knowfreq 45%

basics

~20 s

Setting broker.rack tags each broker with a failure-domain label (e.g. an AZ). Kafka then spreads a partition's replicas across distinct racks so one rack failure doesn't take all replicas down. It improves availability but doesn't change which replica is leader.

open as a page

What does num.recovery.threads.per.data.dir control, and when does increasing it help?

level: seniorimportance: should knowfreq 40%

basics

~20 s

It sets how many threads, per log directory, Kafka uses to load and recover partition logs at startup and to flush them at shutdown. Raising it speeds up startup recovery on brokers with many partitions, especially after an unclean shutdown.

open as a page

Walk through the on-disk layout under a broker's log.dirs: what directories and files exist for a partition, and how do you inspect a segment?

level: seniorimportance: should knowfreq 35%

basics

~10 s

Under each path in log.dirs there's one directory per partition named <topic>-<partition> (e.g. orders-3). Inside are the segment files (.log/.index/.timeindex), a leader-epoch-checkpoint, and partition metadata. You inspect a .log with kafka-dump-log.sh.

open as a page

Explain the head vs tail of a compacted log and why the active segment is never compacted.

level: seniorimportance: should knowfreq 35%

basics

~20 s

The tail is the older, already-compacted part of the log (at most one record per key). The head is newer records appended since the last clean, which may still contain duplicates. The active segment — the one being written to — is always part of the head and is never compacted.

open as a page

Mechanically, how does a single compaction pass deduplicate keys, and what is the offset map?

level: seniorimportance: should knowfreq 40%

basics

~20 s

The cleaner makes two passes over the dirty head. First it builds an in-memory offset map: key -> highest offset seen. Then it recopies the log, keeping each record only if its offset equals the map's value for that key, and merges results into new compacted segments.

open as a page

How does Kafka use memory-mapped files (mmap) for its offset and time indexes, and how do consumers find a record by offset?

level: seniorimportance: should knowfreq 40%

basics

~20 s

Each log segment has sparse index files (.index for offsets, .timeindex for timestamps) that Kafka memory-maps (mmap) so lookups happen in RAM. To find an offset, Kafka binary-searches the sparse index to get a nearby file position, then scans forward in the .log file.

open as a page

Explain watcher keys and the checkAndComplete mechanism: how does a parked operation get woken, and why can an operation register under multiple keys?

level: seniorimportance: should knowfreq 32%

basics

~20 s

Each delayed operation is filed in the purgatory under one or more watcher keys (typically TopicPartition). When an event affecting a key occurs — usually an ISR/high-watermark advance — the broker calls checkAndComplete on that key, which re-runs tryComplete on every operation watching it. Multiple keys let one request span several partitions.

open as a page

What does queued.max.requests control, and what happens when the request queue fills up? Explain the backpressure mechanism.

level: seniorimportance: should knowfreq 40%

basics

~20 s

queued.max.requests (default 500) is the capacity of the shared request queue between network threads and request-handler threads. When it's full, network threads stop reading new requests off sockets, which slows clients down — that is the backpressure.

open as a page

Walk through exactly how Kafka resolves a fetch request for a specific offset using the sparse .index file.

level: seniorimportance: should knowfreq 35%

basics

~20 s

Kafka first picks the right segment (the one whose base offset is the largest <= target). It binary-searches that segment's .index for the entry with the largest offset <= target, seeks to that byte position in the .log, then scans record batches forward until it reaches the requested offset.

open as a page

How does Kafka answer an offset-by-timestamp query (e.g. consumer.offsetsForTimes), and what role does the .timeindex play and what are its caveats?

level: seniorimportance: should knowfreq 30%

basics

~20 s

The consumer sends a ListOffsets request with the target timestamp. The broker uses each segment's .timeindex (timestamp -> offset) to binary-search for the first offset whose timestamp is >= the target, then refines via the .index. It returns that offset and its timestamp.

open as a page

What do flush.messages and flush.ms control, and why does Kafka discourage relying on broker fsync for durability?

level: principalimportance: should knowfreq 45%

basics

~20 s

flush.messages and flush.ms force the broker to fsync log data to disk after N messages or T milliseconds. By default both are effectively unset, so Kafka lets the OS flush dirty pages lazily and relies on replication, not local fsync, for durability — forcing frequent fsync hurts throughput.

open as a page

Tiered Storage relies on an external object store and a metadata store. As a principal engineer, what consistency and durability invariants must hold, and where can things go wrong?

level: principalimportance: should knowfreq 35%

basics

~20 s

A segment must be durably in the object store before its metadata is marked finished/readable; metadata must be the single source of truth for what's readable. Risks include orphaned objects, eventual-consistency reads, metadata/data divergence, and double-offload across leader changes.

open as a page

Why does Kafka use a hierarchical timing wheel for delayed-operation timeouts instead of a java.util.concurrent DelayQueue or a priority queue?

level: principalimportance: nice to knowfreq 22%

basics

~20 s

A timing wheel adds and cancels a timeout in O(1), versus O(log n) for a heap/DelayQueue. Since brokers create and (usually) cancel a huge number of short-lived timeouts that complete early via events, O(1) insert/cancel matters far more than precise ordering.

open as a page

What is the .txnindex file, what does it record, and how do consumers use it to honor read_committed isolation?

level: principalimportance: nice to knowfreq 18%

basics

~20 s

The .txnindex is a per-segment file listing aborted transactions as ranges (producer id, first offset, last stable offset). With read_committed isolation, the broker sends this list so consumers can filter out records from aborted transactions and not read past the last stable offset.

open as a page

showing 31–45 of 45