Interview Questions β Core Kafka Concepts
π― These are the most commonly asked Kafka interview questions, grouped by topic. Each answer is designed to be comprehensive enough to impress, but concise enough to deliver in an interview.
OrderPlaced event to Kafka in <5ms. Independent consumer groups read at their own speed.Architecture
Q1: Explain Kafka's architecture in 2 minutes.
Apache Kafka is a distributed event streaming platform built around an append-only, partitioned log. Producers write messages to topics, which are divided into partitions. Partitions are distributed across multiple brokers in a cluster, with each partition having one leader and zero or more follower replicas. One broker is elected as the Controller, responsible for partition leader elections.
Consumers read from partitions in consumer groups β each partition assigned to exactly one consumer within a group. Kafka retains messages for a configurable time regardless of consumption, enabling replay and multiple independent consumer groups.
In modern Kafka (3.3+), KRaft replaces ZooKeeper for metadata management, embedding a Raft consensus quorum directly into brokers.
Q2: How does Kafka achieve high throughput?
Several design decisions work together: (1) Sequential disk I/O β Kafka appends to the end of log segment files, which is much faster than random writes. (2) OS page cache β Kafka relies on the OS to cache recently written data in memory, avoiding redundant reads from disk. (3) Zero-copy transfer (
sendfilesyscall) β data is transferred directly from page cache to network socket without copying through user space. (4) Batching β producers batch multiple messages; brokers serve batches to consumers. (5) Compression β reduces network and disk I/O.
Q3: What is the role of ZooKeeper in Kafka (legacy) and why is KRaft replacing it?
In legacy Kafka, ZooKeeper stored cluster metadata: broker registry, topic/partition assignments, controller election, and (in old clients) consumer group offsets. This added operational complexity β two distributed systems to manage, monitor, and version-align.
KRaft (Kafka Raft) replaces ZooKeeper by embedding metadata management into Kafka itself using a Raft consensus algorithm. The metadata log is stored in a special
__cluster_metadatatopic managed by a quorum of controller nodes. Benefits: single system to operate, millisecond controller failover (vs tens of seconds), support for millions of partitions per cluster.
Q4: What is the difference between a partition leader and a follower?
The leader handles all reads and writes for a partition. Producers send messages to the leader; consumers fetch from the leader (unless follower fetching is enabled). Followers exist solely to replicate data from the leader, staying in sync to be ready for failover. When the leader fails, the Controller elects a new leader from the ISR (In-Sync Replicas).
Q5: What is the High Watermark and why does it matter?
The High Watermark (HW) is the highest offset that has been replicated to all In-Sync Replicas. Consumers can only read messages up to the HW β they cannot see messages that haven't been fully replicated yet. This prevents consumers from reading data that might be lost if the leader crashes before replication completes. The HW advances as followers acknowledge fetching messages from the leader.
Topics & Partitions
Q6: How do you decide how many partitions to create for a topic?
Consider three factors: (1) Throughput: if a single partition can handle 50 MB/s and you need 200 MB/s, you need at least 4 partitions. (2) Consumer parallelism: the number of partitions limits how many consumers in a group can run in parallel β 6 partitions β max 6 active consumers. (3) Operational overhead: each partition costs memory, file handles, and replication overhead. A common starting point for moderate workloads is 6β12 partitions. Partitions can be increased later (never decreased), so starting conservatively is reasonable.
Q7: What is the difference between delete and compact cleanup policies?
deleteremoves messages based on time (retention.ms) or size (retention.bytes) β oldest segments are deleted. Once deleted, messages are gone.compactretains only the latest value per key indefinitely, deleting records with older versions of the same key. A null-value record (tombstone) signals deletion of a key. Usedeletefor time-series/event streams; usecompactfor state/changelog topics where only current state matters.
Q8: Can you decrease partition count? Why or why not?
No. Decreasing partitions is not supported. The
key β partitionassignment is based onhash(key) % numPartitions. Decreasing partitions would remap keys to different partitions, breaking ordering guarantees and making it impossible to find historical messages for a given key. The only option is to create a new topic with fewer partitions and migrate data.
Q9: What is partition skew and how do you mitigate it?
Partition skew occurs when traffic is unevenly distributed β some partitions get far more messages than others. It's caused by low-cardinality keys (e.g., using
status="ACTIVE"as a key when 90% of messages are active) or viral/hot keys. Mitigation strategies: (1) Use high-cardinality keys (UUIDs, entity IDs). (2) Implement a custom partitioner with special-case logic for hot keys. (3) Add a random salt suffix to hot keys and aggregate in a separate stage. (4) Create a dedicated high-throughput topic for hot entities.
Replication & Durability
Q10: What is the ISR and what happens when a replica falls out of it?
The ISR (In-Sync Replicas) is the set of replicas that are fully caught up with the leader within
replica.lag.time.max.ms. When a follower falls behind β typically due to network issues or broker slowness β it's removed from the ISR. The leader no longer waits for it when advancing the High Watermark. When the follower catches up, it's added back. If the ISR shrinks belowmin.insync.replicas, writes withacks=allare rejected withNotEnoughReplicasException.
Q11: What configuration gives the strongest durability guarantee?
The "zero data loss" combination:
acks=allβ wait for all ISR membersmin.insync.replicas=2β require at least 2 replicas in ISRreplication.factor=3β maintain 3 copiesunclean.leader.election.enable=falseβ never elect out-of-sync replicasenable.idempotence=trueβ prevent duplicate writes on retryThis setup tolerates the loss of 1 broker while maintaining write availability, and never loses an acknowledged message.
Q12: What is unclean leader election? When would you enable it?
Unclean leader election allows an out-of-sync replica (not in ISR) to become the new leader when all ISR members are unavailable. This restores availability at the cost of data loss β the new leader may be missing recently produced messages. It's disabled by default. You might enable it in scenarios where availability is strictly more important than data integrity, such as log aggregation pipelines where losing a few recent log lines is acceptable vs a topic being offline.
Operations
Q13: How would you handle a scenario where consumer lag is growing?
Growing lag means consumers are slower than producers. Diagnose first: (1) Check
max.poll.recordsβ consuming more per poll batch may help. (2) Check processing time per record β if slow, optimize or add parallelism. (3) Check consumer count β add more consumers (up to partition count). (4) Check if a partition is stuck on an error loop consuming retries. (5) Checkmax.poll.interval.msβ if processing is slow, increase it to prevent premature rebalances. If genuinely under-resourced, scale consumer instances horizontally.
Q14: How do you replay messages in Kafka?
Options: (1) Reset consumer group offset β use
kafka-consumer-groups.sh --reset-offsets --to-earliest(stop consumers first). (2) Seek to specific offset or timestamp in application code usingconsumer.seek()orconsumer.offsetsForTimes(). (3) Use a separate consumer group β create a new group that starts from earliest; doesn't disturb existing groups. (4) For full topic replay, mirror the topic to a replay topic and reprocess. Kafka's message retention (default 7 days) makes all of these possible.
Q15: What is rack-aware replica assignment?
Kafka can spread replicas across different racks (physical or availability zones) to ensure that a single rack failure doesn't take a partition offline. Configure
broker.rack=us-east-1aon each broker. When creating topics, Kafka assigns replicas to brokers in different racks. This provides rack-level fault tolerance in addition to broker-level fault tolerance.
Senior Deep-Dive: 6 Critical Kafka System Design Scenarios
Scenario 1: What architectural problem does Kafka solve in an E-Commerce Checkout?
The Problem: In a synchronous architecture, when a user clicks "Place Order", the Checkout Service must call the Payment Service, Inventory Service, Fraud Detection Service, and Notification Service via HTTP/gRPC. If any single service is slow (e.g. Fraud takes 2 seconds) or down (Inventory 500 error), the entire user checkout fails or times out.
The Kafka Solution: The Checkout Service writes a single
OrderPlacedevent to Kafka and returns HTTP 202 / 200 to the user in <5ms. Payment, Inventory, and Notification services each operate in their own independent consumer groups. If the Notification service crashes for 2 hours, it simply resumes from its last committed offset with zero data loss and zero impact on active checkouts.
Scenario 2: How does Kafka organize data on disk for massive throughput?
Kafka organizes data into Topics, which are split into physical Partitions, and partitions are stored as Segment files (
.log,.index,.timeindex) on disk.
- Sequential Append-Only Writes: Kafka never updates data in place; new records are appended sequentially to the active segment file. On modern NVMe/SSDs and spinning disks, sequential writes achieve near memory-bus throughput (600+ MB/s).
- Linux OS Page Cache: Kafka does not cache messages in the JVM heap (avoiding GC overhead). It relies entirely on the OS Page Cache.
- Zero-Copy Network Transfers (
sendfile): When consumers fetch data, the Linux kernel copies bytes directly from the Page Cache to the network socket descriptor without transferring data through user space.
Scenario 3: What is an Offset, and what happens if a consumer crashes?
An Offset is a monotonically increasing 64-bit integer assigned to each record within a partition. It represents the logical position of a message.
When a consumer processes messages, it periodically commits its current offset to the internal
__consumer_offsetstopic.
- During Crash / Failover: If Consumer A crashes mid-processing, Kafka detects heartbeat timeout, triggers a Consumer Group Rebalance, and assigns the partition to Consumer B. Consumer B reads
__consumer_offsetsand starts reading from the last committed offset.- At-Least-Once Risk: If Consumer A processed the record but crashed before committing offset, Consumer B will re-process that record. Hence, all downstream business operations must be idempotent.
Scenario 4: How does Key-Based Partitioning guarantee ordering?
Kafka guarantees message order only within a single partition, never across multiple partitions of a topic.
- When a producer sends a record with a non-null key (e.g.,
customerIdororderId), Kafka calculates:- Because this hashing algorithm is deterministic, every single event for customer
cust_88will ALWAYS route to the exact same partition in the exact sequence they were published.- Gotcha: If you dynamically increase the partition count on a live topic, the modulo divisor changes, re-mapping existing keys to new partitions and breaking ordering for in-flight streams.
Scenario 5: Why does a Consumer Group have an upper limit on active workers?
The 1:1 Partition Assignment Rule: Within a single consumer group, a partition can be consumed by at most ONE active consumer instance at any given time.
- If topic
ordershas 12 partitions:
- 4 consumers each consumes 3 partitions.
- 12 consumers each consumes 1 partition (maximum useful parallelism).
- 20 consumers 12 active consumers, while 8 instances sit completely idle as hot standbys.
- To increase consumer parallelism beyond 12 workers, you must increase the partition count of the topic.
Scenario 6: Does Kafka TRULY guarantee "Exactly-Once Processing"?
The Interview Trap Answer: "Yes, Kafka has
processing.guarantee=exactly_once_v2." (Junior answer).The Senior Architect Answer: "Kafka's Exactly-Once Semantics (EOS) only applies within the Kafka ecosystem (reading from Kafka, processing in Kafka Streams / Flink, and writing to another Kafka topic). Inside that boundary, Kafka uses idempotent producer sequence numbers and 2-Phase Commit transaction markers across
__transaction_stateand output partitions.""However, the moment a consumer writes to an external system (e.g. updating a PostgreSQL database, charging a Stripe credit card, or sending an SMS), exactly-once is mathematically impossible via Kafka alone. If the consumer updates the database but crashes before committing its offset to Kafka, the message will be redelivered. Therefore, end-to-end exactly-once requires Application-Level Deduplication (idempotency keys + unique DB constraints or Redis
SETNX)."
Scenario 7: An interviewer asks: "What happens if Kafka goes down?" How do you respond?
The Principal Architect Response: "I gently clarify the failure boundary. Kafka is architected as an always-available, horizontally distributed commit log. In an enterprise production deployment across multiple Availability Zones with a replication factor of 3 () and KRaft metadata quorum, an entire cluster-wide outage is exceedingly rare unless there is a catastrophic global cloud provider networking failure."
"Instead, what is far more realistic and critical to design for in an interview are localized component failures:
- Broker Leader Hardware Failure: The KRaft controller instantly detects the loss via heartbeat timeout and promotes an In-Sync Replica (ISR) follower to leader with zero data loss (
acks=all+min.insync.replicas=2).- Consumer Instance Crash: The Group Coordinator detects
session.timeout.msexpiration and triggers a rebalance to reassign orphaned partitions to surviving consumers.- Producer Network Partition: The producer buffers records in memory (
RecordAccumulator) and retries automatically with idempotent sequence numbers once connectivity is restored.*
Scenario 8: How do you handle viral events and Hot Partitions in Kafka (e.g. Nike LeBron James ad in Ad Click Aggregator)?
In an Ad Click Aggregator where click events are partitioned by
ad_id, a viral campaign (e.g., Nike LeBron James ad during the NBA Finals) causes 90% of global traffic to map to a single partition (abs(murmur2("nike_ad")) % num_partitions). This overwhelms the broker's disk I/O and creates an intractable consumer lag bottleneck.The 4 Production Solutions:
- Omit the Key (Default Sticky Partitioner): If ordering is not strictly required (e.g., calculating aggregate click sums where addition is associative and commutative), remove the key. Kafka's sticky partitioner batches records and distributes them evenly across all partitions.
- Random Salting with Two-Stage Aggregation: Append a random integer suffix (
ad_id + "#salt=" + random(10)) to spread the viral ad across 10 partitions. Downstream stream processing workers (Flink or Kafka Streams) run a two-stage aggregation: local tumbling window by salted key, followed by a global merge by originalad_id.- Compound Key Partitioning: Combine the entity ID with an independent variable that distributes evenly, such as
ad_id + "#" + user_regionorad_id + "#" + (user_id % 16).- Producer Backpressure: The ingestion API monitors partition consumer lag and throttles incoming client requests via HTTP 429 when lag breaches safety thresholds.
Scenario 9: Can Kafka store videos or large images? (The Claim Check Pattern)
Never store large binary blobs directly in Kafka. Kafka is engineered for small event payloads (1 KB to 100 KB). Large payloads (> 1 MB) pollute the OS Page Cache, cause massive JVM garbage collection pauses, and degrade throughput.
Instead, implement the Claim Check Pattern:
- The client uploads the raw video or image directly to object storage (Amazon S3 or Google Cloud Storage) via a presigned URL.
- The API service publishes a tiny metadata record (< 1 KB) to Kafka containing the S3 URI pointer and processing parameters.
- Downstream worker pools (e.g., transcoding workers in a YouTube system design) pull the pointer from Kafka and stream chunks directly from S3.
