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 pageshowhide
explore
- Broker and Cluster Anatomy5 questions
- Partition Commit Log and Segments5 questions
- Offset and Time Indexes5 questions
- Page Cache and Zero-Copy I/O5 questions
- Request Processing Pipeline5 questions
- Request Purgatory and Delayed Operations5 questions
- Log Cleaner and Compaction Internals5 questions
- Broker Memory, Threads and Disk Layout5 questions
- Tiered Storage Internals5 questions
questions
page 2 of 2Walk through how a closed log segment gets offloaded to remote storage. What does the broker upload, and when can it delete the local copy?
basics
~20 sWhen 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.
How does broker.rack enable rack awareness, and what does it (and doesn't it) guarantee for replica placement and reads?
basics
~20 sSetting 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.
What does num.recovery.threads.per.data.dir control, and when does increasing it help?
basics
~20 sIt 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.
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?
basics
~10 sUnder 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.
Explain the head vs tail of a compacted log and why the active segment is never compacted.
basics
~20 sThe 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.
Mechanically, how does a single compaction pass deduplicate keys, and what is the offset map?
basics
~20 sThe 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.
How does Kafka use memory-mapped files (mmap) for its offset and time indexes, and how do consumers find a record by offset?
basics
~20 sEach 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.
Explain watcher keys and the checkAndComplete mechanism: how does a parked operation get woken, and why can an operation register under multiple keys?
basics
~20 sEach 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.
What does queued.max.requests control, and what happens when the request queue fills up? Explain the backpressure mechanism.
basics
~20 squeued.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.
Walk through exactly how Kafka resolves a fetch request for a specific offset using the sparse .index file.
basics
~20 sKafka 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.
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?
basics
~20 sThe 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.
What do flush.messages and flush.ms control, and why does Kafka discourage relying on broker fsync for durability?
basics
~20 sflush.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.
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?
basics
~20 sA 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.
Why does Kafka use a hierarchical timing wheel for delayed-operation timeouts instead of a java.util.concurrent DelayQueue or a priority queue?
basics
~20 sA 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.
What is the .txnindex file, what does it record, and how do consumers use it to honor read_committed isolation?
basics
~20 sThe .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.
showing 31–45 of 45