skip to content

Apache Kafka

3 roadmaps835 questionsupdated

Apache Kafka end to end: the partitioned commit log, producers and consumer groups, replication, Streams and Connect, schemas, security, and multi-cluster operations. Kafka shows up in most backend and data-engineering interviews, so the whole surface is fair game.

on this pageshow

guide

overview

~1 min

Apache Kafka is a distributed, partitioned commit log: producers append records to topics, brokers keep them for a configured time, and any number of consumer groups read them back at their own pace. Interviewers lean on it because it sits in the middle of so many backend and data systems, and because its guarantees are precise enough to probe. A good Kafka answer names the setting, the failure it guards against, and the price paid for it. The hub splits along the path a record takes. [Topics, partitions and log storage](/topics/cloud-kafka-topics-partitions) is the data model everything else refers to. [Producers](/topics/cloud-kafka-producers) and [consumers and consumer groups](/topics/cloud-kafka-consumers) are the write and read paths, where most loss, duplicate and lag stories begin. [Replication and durability](/topics/cloud-kafka-replication-durability) and [delivery semantics and transactions](/topics/cloud-kafka-delivery-semantics) explain what "the write succeeded" actually promises. Beneath the clients sit [architecture and internals](/topics/cloud-kafka-architecture) and [cluster coordination](/topics/cloud-kafka-kraft). Built on the core are [Kafka Streams](/topics/cloud-kafka-streams), [Kafka Connect](/topics/cloud-kafka-connect) and [schemas and serialization](/topics/cloud-kafka-schema-serialization). The operator's side covers [security](/topics/cloud-kafka-security), [operations](/topics/cloud-kafka-operations), [monitoring and performance](/topics/cloud-kafka-monitoring-performance) and [multi-cluster replication](/topics/cloud-kafka-multi-cluster). [Client patterns](/topics/cloud-kafka-client-patterns) is where application developers meet all of this in code, and [ecosystem and platform choices](/topics/cloud-kafka-ecosystem) asks when Kafka is the wrong tool. Junior rounds stay on the model: what a partition is, what an offset means, what a consumer group shares. Senior and staff rounds become trade-off conversations — durability against availability, ordering against parallelism, exactly-once against its cost — and incident stories where you are expected to know which signal to check first. Learn the log and its partitions before anything else, then the producer and consumer settings that decide durability and duplicates. Replication and transactions make sense only once those are solid; Streams, Connect and operations build on all three.

primer

### The log is the product Kafka stores data in **partitions**: ordered, append-only sequences of records, each addressed by an offset. Topics are named groups of partitions. Reading removes nothing — records leave when retention or compaction says so — which is why many independent consumers can share the same data, and why replaying history is routine rather than a recovery trick. Most design answers reduce to one question: what goes into a partition, and what does that buy or cost. ### Partitions are the unit of everything - **Ordering** holds inside one partition and nowhere else, so the record key, which picks the partition, is a design decision rather than a detail. - **Parallelism** for a group of consumers is capped by partition count; members beyond it sit idle. - **Replication, leadership and placement** are all tracked per partition, so a "cluster problem" is usually a set of partition problems. Raising the partition count later remaps keys, so the number is chosen early and defended in design rounds. ### Durability is a dial, not a switch How safe a write is depends on three settings chosen together: the number of copies, the number that must stay in sync before a write is accepted, and how many acknowledgements the producer demands. The **ISR** and the **high watermark** decide what readers may see and what survives a leader change. Interviewers want you to turn the dial both ways and name what each notch of safety costs in latency or availability. ### Consumers own their position A consumer group divides a topic's partitions among its members, and each member records progress by committing offsets. Where the commit sits relative to the processing decides whether a crash loses records or repeats them. Rebalancing — members joining, leaving or timing out — is when those choices get tested, and the source of most consumer incidents. ### Topics outlive their writers Because records stay and get replayed, the format on a topic is a long-lived contract read by services that did not exist when it was written. Evolving it compatibly matters as much as versioning an API, which is what serializers, schemas and the registry are for. ### Exactly-once has a boundary Idempotent producers and transactions make a read-process-write loop atomic across output records and consumer offsets inside one cluster. Outside that loop — a database, an HTTP call, a second cluster — the guarantee stops and application-level idempotence takes over. A strong answer draws the boundary before explaining the machinery. ### The broker is simple on purpose Brokers store files and serve bytes; batching, compression and partition choice mostly happen in the clients. Sequential appends, the OS page cache and zero-copy transfer explain the throughput, and why tuning a broker looks unlike tuning a database.

Topic
A named stream of records, split into one or more partitions. It is a logical grouping; storage, ordering and replication all happen per partition.
Partition
One ordered, append-only log within a topic, stored on a leader broker and copied to followers. The unit of ordering, parallelism and replication.
Offset
The sequential position of a record within its partition. Consumers use it as a bookmark; it never identifies a record across partitions.
Record key
An optional value attached to a record. The default partitioner hashes it to choose a partition, so records sharing a key stay together and in order.
Log segment
One file of a partition's log on disk. Retention deletes whole segments and compaction rewrites closed ones; only the newest, active segment receives appends.
Consumer group
A set of consumers sharing a group id that split a topic's partitions between them, so each partition is read by one member at a time.
Rebalance
The redistribution of partitions among a group's members, triggered when members join, leave, time out or subscriptions change.
Committed offset
The position a consumer group has stored for a partition, used to resume after a restart or reassignment.
Consumer lag
The distance between the newest record in a partition and the group's committed position; the first signal that consumers cannot keep up.
In-sync replicas (ISR)
The replicas of a partition, leader included, that are caught up with the leader closely enough to be trusted with acknowledged data and leadership.
High watermark
The offset up to which records are replicated to the whole ISR. Consumers read only below it, so they do not see data a clean failover could truncate.
acks
The producer setting that decides how many broker confirmations a send waits for: none, the leader only, or every in-sync replica.
min.insync.replicas
The smallest ISR size a partition needs to accept a write sent with acks=all. Below it, such writes are rejected rather than under-replicated.
Idempotent producer
A producer whose retries cannot create duplicates within a partition, because brokers track a producer id and per-partition sequence numbers.
Transaction
A group of writes across partitions, optionally including consumer offsets, that becomes visible atomically to read_committed consumers or not at all.
Log compaction
A cleanup policy that keeps at least the latest record per key instead of deleting by age; a record with a null value, a tombstone, eventually removes the key.
KRaft
Kafka's built-in Raft-based metadata quorum. A small set of controllers holds cluster metadata in a replicated log, replacing the older ZooKeeper dependency.
Schema Registry
A separate service that stores versioned schemas for topic data and enforces compatibility rules, so producers and consumers agree on record shape.

Follow one record through the system. A producer serializes key and value, picks a partition from the key, and batches records per partition before sending them to that partition's leader broker. The leader appends to its log, followers fetch from it, and once enough in-sync replicas hold the batch the high watermark moves forward and the producer receives its acknowledgement. Each consumer in a group owns some partitions, fetches from their leaders (or a nearby follower in multi-region setups), processes, and commits offsets back into Kafka. Metadata — which broker leads which partition, who is in the ISR — lives in the KRaft controllers' metadata log, which brokers follow and clients query through any broker. The rest of the hub is layered on this path rather than beside it: - **Transactions** add a coordinator and commit markers, so output records and consumer offsets become visible together. - **Kafka Streams** is a client library; its state stores are backed by compacted changelog topics, and its parallelism follows the input partitions. - **Kafka Connect** runs producers and consumers for you inside workers, driven by configuration; in distributed mode it keeps its own offsets and configs in topics. - **Schema Registry** sits outside the brokers; the serializer writes a schema id into each record, and to a plain broker the payload is still bytes. - **MirrorMaker 2** is built on Connect, consuming from one cluster and producing into another. The operator's sections attach to the same objects. ACLs and quotas name topics, groups and principals; lag, under-replication and request latency are read per partition or per broker; cross-cluster replication copies partitions, and with MirrorMaker 2 it must also translate consumer offsets, because a mirrored partition's offsets need not match the source. The durability dial spans three owners — the topic, the producer and the consumer — which is why an answer that sets only one of them falls short: ```properties # topic / broker: three copies, two must be in sync to accept a write default.replication.factor=3 min.insync.replicas=2 # producer: wait for the whole ISR, retries cannot duplicate acks=all enable.idempotence=true # consumer: commit after the work, skip aborted transactional data enable.auto.commit=false isolation.level=read_committed ``` No line is enough alone. `acks=all` against a minimum ISR of one can still acknowledge data held by a single broker, and committing after processing still repeats work after a crash unless the processing itself is idempotent.

  1. Topics, Partitions and Log Storage →

    The data model every other section assumes: partitions, keys, ordering, offsets and retention.

  2. Producers →

    The write path, where acks, batching and idempotence decide whether data is lost or duplicated.

  3. Consumers and Consumer Groups →

    The read path: groups, rebalancing and offset commits, where most production incidents first show up.

  4. Replication and Durability →

    What an acknowledged write really survives once brokers fail, and the durability-versus-availability trade.

  5. Delivery Semantics and Transactions →

    Delivery guarantees and transactions make sense only after producers, consumers and replication are clear.

  6. Client Development and Integration Patterns →

    Where the pieces meet in application code: retries, dead-letter topics, outbox and testing.

  • Claiming Kafka orders a topic: it orders a partition, and changing a key or the partition count quietly moves records to a different partition.

  • Saying acks=all means every replica: it means every replica currently in sync, which can shrink to the leader alone unless min.insync.replicas prevents it.

  • Presenting exactly-once as end to end: it covers topics and offsets inside one cluster, and a database write or HTTP call in the loop needs its own idempotence.

  • Committing offsets before processing finishes, or leaning on auto-commit without knowing when it fires — the usual root cause of both lost and repeated records.

  • Answering growing lag with more consumers when the group already has one member per partition; extra members stay idle and only add rebalance churn.

  • Blaming the broker for a rebalance storm caused by slow work inside the poll loop — see Max Poll Interval and Long-Processing Pitfalls.

  • Recommending unclean leader election for availability without saying it can discard writes that producers were already told had succeeded.

  • Describing ZooKeeper as a current dependency: Kafka 4.x runs KRaft only, so name the version whenever an answer relies on ZooKeeper-era behaviour.

This guide assumes Apache Kafka 4.x, which runs only in KRaft mode. Several questions still hinge on what changed along the way, because plenty of production clusters and client libraries are older: - **The 3.x line** changed the Java producer defaults to `acks=all` with idempotence enabled, so "what are the default durability settings" has two answers depending on client version. - **3.3** marked KRaft production-ready for new clusters. - **3.9** is the last release line that supports ZooKeeper and the bridge for migrating a ZooKeeper cluster to KRaft; it also added dynamic controller quorum membership. - **4.0** removed ZooKeeper entirely and made the new consumer group protocol generally available, moving assignment to the broker-side coordinator and removing the group-wide synchronization barrier from rebalances. When an answer depends on version — producer defaults, the rebalance protocol, how the controller is chosen — say which version you are describing. Interviewers who run older clusters will ask about ZooKeeper-era behaviour, and naming the era is part of a good answer.

Interviewers expect you to place Kafka, not only operate it. As a log, it competes most directly with other log-based systems — Apache Pulsar, and managed or self-hosted services that speak the Kafka wire protocol. It is also weighed against traditional brokers such as classic RabbitMQ queues, which drop a message after acknowledgement and handle per-message routing, priorities and redelivery well. A defensible choice names the workload: replayable history, high fan-out and ordered per-key streams favour a log; per-message work queues often favour a queue. Around the core sit layers you will be asked to choose among. Kafka Streams and ksqlDB process data inside the Kafka world; Apache Flink and Spark Structured Streaming run as their own clusters and fit when state, windowing or sources outgrow a library. Kafka Connect, with Debezium for change data capture, replaces most hand-written ingestion producers. Schema Registry, from Confluent, is the common answer for typed contracts on topics. Managed offerings remove broker operations but not partitioning, keying or consumer design — those questions follow you to every provider. [Ecosystem and Platform Choices](/topics/cloud-kafka-ecosystem) covers the comparisons in depth.

explore

report an issue with this guide →

questions

835 · 16 sections

What is the purpose of bootstrap.servers, and why don't you need to list every broker?

level: juniorimportance: must knowfreq 80%
basics
~20 s

bootstrap.servers is the initial list of broker host:port pairs a client contacts to discover the full cluster. After the first metadata fetch the client learns all brokers, so the list only needs a few entries for redundancy.

open as a page

What is a Kafka broker, and what role does broker.id play in a cluster?

level: juniorimportance: must knowfreq 75%
basics
~10 s

A broker is a single Kafka server that stores topic partition data and serves produce/fetch requests. broker.id is its unique numeric identifier within the cluster; no two brokers may share the same id.

open as a page

What is a Kafka partition's commit log, and why is it split into segments on disk?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Each partition is an append-only log: records are only added to the end, never changed in place. Kafka splits that log into fixed-size files called segments so old data can be deleted or compacted one whole file at a time instead of editing one giant file.

open as a page

What is the Kafka log cleaner, and how do you enable it for a topic?

level: juniorimportance: must knowfreq 60%
basics
~20 s

The log cleaner is a background process that runs compaction: it scans a topic's log and keeps only the latest record per key, deleting older duplicates. You enable it by setting the topic config cleanup.policy=compact.

open as a page

Why does Kafka rely on the operating system page cache instead of maintaining its own in-process (JVM heap) record cache?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Kafka writes data to files and lets the OS keep recently used file pages in RAM (the page cache). It avoids a JVM heap cache to dodge GC pressure, double-buffering, and to reuse the OS cache that survives broker restarts.

open as a page

In KRaft mode, what is the active controller and what role does it play in managing cluster metadata?

level: juniorimportance: must knowfreq 70%
basics
~20 s

The active controller is the single elected controller node that is the only one allowed to write cluster metadata (topics, partitions, broker registrations) by appending records to a replicated log. Other nodes follow and replay that log.

open as a page

What is KRaft in Apache Kafka, and what problem was it introduced to solve?

level: juniorimportance: must knowfreq 80%
basics
~10 s

KRaft (Kafka Raft) is Kafka's built-in consensus protocol that replaces Apache ZooKeeper for storing cluster metadata. It lets Kafka manage its own metadata internally, so you no longer run a separate ZooKeeper cluster.

open as a page

In a KRaft Kafka cluster, how does a broker tell the controller it is still alive, and what RPC is involved?

level: juniorimportance: must knowfreq 55%
basics
~10 s

Each broker periodically sends a BrokerHeartbeat request to the active controller. The interval is set by broker.heartbeat.interval.ms (default 2000 ms). If heartbeats stop, the controller eventually fences the broker.

open as a page

What is the difference between the static controller.quorum.voters config and the dynamic quorum introduced by KIP-853?

level: juniorimportance: must knowfreq 60%
basics
~10 s

Static quorums list every controller in controller.quorum.voters on all nodes, fixed at startup. KIP-853 dynamic quorums (KRaft) let you add or remove controllers at runtime with AddRaftVoter/RemoveRaftVoter instead, using controller.quorum.bootstrap.servers.

open as a page

In a KRaft cluster, why does Kafka periodically take metadata snapshots of the __cluster_metadata log?

level: juniorimportance: must knowfreq 62%
basics
~20 s

The metadata log grows forever as records are appended. Snapshots capture the current state at a point in time so old log records can be deleted, keeping the log small and making new nodes catch up faster.

open as a page

When you send a Kafka record with a non-null key, how does the producer decide which partition it goes to, and why does that matter?

level: juniorimportance: must knowfreq 78%
basics
~10 s

The producer hashes the key and takes hash % numberOfPartitions. The same key always lands in the same partition, so all records for that key stay ordered together on one partition.

open as a page

What does cleanup.policy=compact do to a Kafka topic, and how is it different from the default delete policy?

level: juniorimportance: must knowfreq 70%
basics
~10 s

Compaction keeps at least the latest value for each message key, deleting older values for that key. The default delete policy instead drops whole segments once they exceed retention.ms or retention.bytes, regardless of key.

open as a page

What is an offset in Kafka, and why are offsets meaningful only within a single partition?

level: juniorimportance: must knowfreq 78%
basics
~20 s

An offset is a number that marks a record's position inside one partition. It starts at 0 and increases by 1 for each new record. Offsets are per-partition, so offset 5 in partition 0 and offset 5 in partition 1 are unrelated records.

open as a page

What ordering guarantee does Kafka provide, and what does it NOT guarantee?

level: juniorimportance: must knowfreq 80%
basics
~10 s

Kafka guarantees order only within a single partition: messages are read in the same order they were written. It does NOT guarantee any order across different partitions of a topic.

open as a page

What does `kafka-topics.sh --alter --partitions` allow you to do, and what is the key directionality constraint?

level: juniorimportance: must knowfreq 70%
basics
~10 s

It can only INCREASE the partition count of an existing topic, never decrease it. To shrink, you must create a new topic with fewer partitions and republish the data; Kafka has no in-place shrink.

open as a page

What do the producer acks settings 0, 1, and all (-1) mean, and how do they trade durability against latency?

level: juniorimportance: must knowfreq 85%
basics
~20 s

acks=0 means the producer never waits for acknowledgement (fastest, can lose data). acks=1 waits only for the leader to write the record. acks=all waits for the leader plus all in-sync replicas, giving the strongest durability but the highest latency.

open as a page

What do batch.size and linger.ms control in a Kafka producer, and how do they trade latency for throughput?

level: juniorimportance: must knowfreq 78%
basics
~20 s

batch.size caps how many bytes a producer collects per partition before sending; linger.ms tells it to wait up to that many milliseconds for more records to fill a batch. Bigger batches and more lingering raise throughput but add latency.

open as a page

What is the difference between a retriable and a fatal (non-retriable) exception in the Kafka producer, and how does each affect a send?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Retriable errors (like a leader change or timeout) are transient, so the producer automatically resends. Fatal errors (like an unknown topic, message too large, or auth failure) cannot be fixed by resending, so the send fails immediately.

open as a page

What is Kafka's idempotent producer, and how do you turn it on?

level: juniorimportance: must knowfreq 78%
basics
~10 s

An idempotent producer guarantees that producer retries don't create duplicate messages on a partition. You enable it by setting enable.idempotence=true (the default since Kafka 3.0).

open as a page

How does a Kafka producer decide which partition a record goes to when the record has a non-null key?

level: juniorimportance: must knowfreq 78%
basics
~20 s

For a keyed record, Kafka hashes the key and maps it to a partition. The same key always lands on the same partition (as long as the partition count is unchanged), which keeps records with that key in order.

open as a page

What are key.deserializer and value.deserializer in a Kafka consumer, and why must they pair with the producer's serializers?

level: juniorimportance: must knowfreq 80%
basics
~10 s

Kafka stores message keys and values as raw bytes. key.deserializer and value.deserializer are classes that turn those bytes back into objects. They must match the producer's serializers, or the bytes won't decode correctly.

open as a page

What is the Group Coordinator in Kafka, and how is the coordinator broker for a particular consumer group chosen?

level: juniorimportance: must knowfreq 70%
basics
~10 s

The Group Coordinator is a broker that manages a consumer group: it handles members joining/leaving, triggers rebalances, and stores offsets. The coordinator is the broker that leads the __consumer_offsets partition the group maps to.

open as a page

What is a consumer heartbeat in Kafka, and what does it accomplish?

level: juniorimportance: must knowfreq 70%
basics
~20 s

A heartbeat is a small periodic signal a consumer sends to the group coordinator broker to prove it is alive and still part of its consumer group. If heartbeats stop, the broker assumes the consumer died and reassigns its partitions.

open as a page

How do you inspect consumer lag from the command line, and what do the columns of kafka-consumer-groups.sh --describe mean?

level: juniorimportance: must knowfreq 70%
basics
~10 s

Run kafka-consumer-groups.sh --bootstrap-server host:port --describe --group <group>. It lists each partition with CURRENT-OFFSET (committed), LOG-END-OFFSET (latest), and LAG = the difference, plus the owning consumer.

open as a page

What is consumer lag in Kafka, and how is it calculated for a single partition?

level: juniorimportance: must knowfreq 80%
basics
~10 s

Consumer lag is how far behind a consumer is on a partition: the latest message offset (log-end-offset) minus the offset the consumer has committed. Lag of 0 means fully caught up.

open as a page

What are the three main Kafka settings that work together to control write durability, and what does each do at a high level?

level: juniorimportance: must knowfreq 78%
basics
~20 s

Replication factor (how many copies of each partition), min.insync.replicas (how many copies must be caught up to accept a write), and acks (how many copies the producer waits for: 0, 1, or all). Together they decide how many failures you can survive without losing data.

open as a page

When the broker hosting a partition's leader replica dies, what happens to that partition so producers and consumers can keep working?

level: juniorimportance: must knowfreq 78%
basics
~10 s

Kafka promotes one of the in-sync follower replicas to be the new leader. Clients get an error, refresh metadata, find the new leader's broker, and resume producing/consuming there. No manual action is needed.

open as a page

After a Kafka write returns successfully, where does the data physically live, and why does that distinction matter?

level: juniorimportance: must knowfreq 60%
basics
~20 s

After a successful write the data is in the broker's RAM (the OS page cache), and with acks=all also in the RAM of other replicas. It isn't necessarily on the physical disk yet — the operating system writes it to disk a bit later.

open as a page

What is the difference between the log-end-offset (LEO) and the high watermark (HW) on a Kafka partition leader, and why does it matter to consumers?

level: juniorimportance: must knowfreq 70%
basics
~20 s

LEO is the offset just past the last message written to the log. HW is the highest offset that has been replicated to all in-sync replicas. Consumers can only read up to (below) the HW, so they never see un-replicated records.

open as a page

What is the In-Sync Replicas (ISR) set in Apache Kafka, and how does it relate to the assigned replicas of a partition?

level: juniorimportance: must knowfreq 78%
basics
~20 s

The ISR is the subset of a partition's replicas (leader plus followers) that are fully caught up with the leader's log. Assigned replicas (AR) are all replicas; ISR is the healthy, up-to-date ones eligible to serve and be elected leader.

open as a page

What is the consume-transform-produce pattern in Kafka, and why does plain at-least-once delivery fall short of exactly-once for it?

level: juniorimportance: must knowfreq 70%
basics
~20 s

It's the read-process-write loop: a consumer reads from an input topic, the app transforms records, and a producer writes results to an output topic. At-least-once can produce duplicates or commit offsets without the output, so the output and the offset commit aren't atomic.

open as a page

What is the scope boundary of Kafka's exactly-once semantics (EOS), and why doesn't it automatically extend to external systems?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Kafka EOS only guarantees exactly-once within a single Kafka cluster: a read-process-write where the input topic, output topic, and consumer offsets all live in that same cluster. External databases, APIs, or other clusters are outside the transaction, so it can't cover them.

open as a page

What are the three delivery semantics in Kafka (at-most-once, at-least-once, exactly-once), and what does each guarantee about message delivery?

level: juniorimportance: must knowfreq 80%
basics
~20 s

At-most-once: each message is delivered zero or one time (loss possible, no duplicates). At-least-once: delivered one or more times (duplicates possible, no loss) — Kafka's default. Exactly-once: delivered once and only once (no loss, no duplicates).

open as a page

What does it mean to make a Kafka consumer idempotent, and why does that let you live safely with at-least-once delivery?

level: juniorimportance: must knowfreq 75%
basics
~20 s

An idempotent consumer can process the same message more than once with no extra effect. Kafka's at-least-once delivery can redeliver a record after a crash; if processing is idempotent, those duplicates are harmless, so you get effectively exactly-once results.

open as a page

What does the consumer setting isolation.level do in Kafka, and what are its two possible values?

level: juniorimportance: must knowfreq 65%
basics
~10 s

isolation.level controls whether a consumer sees records from in-progress or aborted transactions. read_uncommitted (default) returns all records; read_committed only returns records from committed transactions, hiding aborted and not-yet-committed ones.

open as a page

In Kafka Streams, what is the difference between groupByKey() and groupBy(), and why does one of them trigger a repartition?

level: juniorimportance: must knowfreq 70%
basics
~20 s

groupByKey() groups by the existing record key with no repartition. groupBy() picks a NEW key, so Streams must repartition (re-shuffle records to topic partitions) so all records with the same new key land on the same task.

open as a page

What is the difference between a KStream and a KTable in the Kafka Streams DSL, and when would you use each?

level: juniorimportance: must knowfreq 85%
basics
~20 s

A KStream is a record stream where every record is an independent event (insert/append). A KTable is a changelog stream where each record is an update keyed by its key — only the latest value per key matters.

open as a page

What does setting processing.guarantee=exactly_once_v2 in a Kafka Streams application actually guarantee, and how do you enable it?

level: juniorimportance: must knowfreq 70%
basics
~20 s

It guarantees each input record affects the application's state and output exactly once, even if the app crashes and restarts — no duplicates and no lost updates. You enable it by setting the config processing.guarantee to exactly_once_v2.

open as a page

What are Interactive Queries in Kafka Streams, and how do you read the value for a key from a local state store?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Interactive Queries let your app read directly from the state stores Kafka Streams already maintains, instead of querying an external database. You call KafkaStreams.store(...) to get a read-only store and look up a key with get().

open as a page

What kinds of joins does Kafka Streams support, and how do they differ at a high level?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Kafka Streams supports stream-stream joins (over a time window), stream-table joins (enrich a record with the latest table value), and table-table joins (combine two changelog tables). GlobalKTable joins let a stream join a fully-replicated table without matching partitioning.

open as a page

What is Change Data Capture (CDC) with Debezium, and why is log-based CDC preferred over query-based polling?

level: juniorimportance: must knowfreq 78%
basics
~20 s

Debezium is a Kafka Connect source connector that reads a database's transaction log (MySQL binlog, Postgres WAL) and streams every row insert, update, and delete to Kafka as events. Log-based CDC catches all changes, including deletes, with low impact on the database.

open as a page

What is a converter in Kafka Connect, and what is the difference between key.converter and value.converter?

level: juniorimportance: must knowfreq 80%
basics
~10 s

A converter serializes/deserializes Connect data to/from bytes on the Kafka topic. key.converter handles the record key; value.converter handles the record value. They are configured independently and can differ.

open as a page

What does errors.tolerance control in Kafka Connect, and what is the difference between the values none and all?

level: juniorimportance: must knowfreq 78%
basics
~10 s

errors.tolerance decides what Connect does when a record fails processing. none (the default) stops the connector task on the first error; all skips the bad record and keeps going.

open as a page

What is Kafka Connect, and what problem does it solve compared to writing your own producer/consumer applications?

level: juniorimportance: must knowfreq 70%
basics
~10 s

Kafka Connect is a framework for streaming data between Kafka and external systems (databases, files, S3) using ready-made connectors, so you configure plugins instead of writing custom producer/consumer code.

open as a page

How does Kafka Connect track offsets for source connectors versus sink connectors? Where is each kind stored?

level: juniorimportance: must knowfreq 75%
basics
~20 s

Source connectors store their own progress (where they read from the external system) in a special Kafka topic called the offset storage topic. Sink connectors just use a normal Kafka consumer group, so Kafka tracks their offsets like any consumer.

open as a page

What is Apache Avro and why is it commonly used as the serialization format for Kafka records?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Avro is a compact binary serialization format that stores data with a separate schema. With Kafka it gives small messages, a typed contract for producers and consumers, and safe schema evolution via a Schema Registry.

open as a page

What are the compatibility modes in Confluent Schema Registry, and at a high level what does each one allow?

level: juniorimportance: must knowfreq 75%
basics
~20 s

Schema Registry has BACKWARD (default), FORWARD, FULL, NONE, and a _TRANSITIVE variant of each. They control what schema changes are allowed: BACKWARD lets new consumers read old data, FORWARD lets old consumers read new data, FULL means both, NONE disables checks.

open as a page

Why is deserializing a Kafka message treated as security-sensitive, and what is the core threat when a consumer deserializes an untrusted payload?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Deserialization turns raw bytes back into objects. A Kafka topic is just bytes from whoever produced them, so a malicious producer can send crafted bytes that exploit the consumer's deserializer — for example triggering code execution or crashing it.

open as a page

What is a "poison pill" record in a Kafka consumer, and why can it stall a consumer group?

level: juniorimportance: must knowfreq 70%
basics
~20 s

A poison pill is a record the consumer cannot deserialize (corrupt or wrong-format bytes). Deserialization happens inside poll(), so it throws every time, the offset never advances, and the consumer is stuck retrying the same record forever.

open as a page

Confluent Schema Registry supports Avro, Protobuf, and JSON Schema. How do you produce Protobuf or JSON Schema messages, and how does the registry know which format a schema is?

level: juniorimportance: must knowfreq 60%
basics
~10 s

Use the format-specific serializer: KafkaProtobufSerializer or KafkaJsonSchemaSerializer instead of KafkaAvroSerializer. Each subject in the registry carries a schemaType (AVRO, PROTOBUF, or JSON) so the registry knows which format the stored schema is.

open as a page

How do you turn on an audit log of allowed and denied authorization decisions in Apache Kafka, and where does that output go?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Kafka's authorizer writes an audit trail via a dedicated logger named kafka.authorizer.logger. Set that log4j logger to DEBUG to see allowed requests too; at INFO you only see denials. Route it to its own file appender.

open as a page

What is a Kafka ACL, and what are the fields that make up a single ACL binding?

level: juniorimportance: must knowfreq 70%
basics
~20 s

An ACL (access control list entry) is a rule saying a principal (user) is allowed or denied a specific operation (like Read or Write) on a resource (like a topic), optionally from a specific host.

open as a page

What is the super.users setting in Kafka, and what happens when a principal is listed in it?

level: juniorimportance: must knowfreq 70%
basics
~10 s

super.users is a broker config listing principals that are granted full access and bypass all ACL checks. Any principal in the list is always authorized for every operation, no ACLs needed.

open as a page

Does Apache Kafka natively encrypt the data it writes to disk (log segments)? If not, how do teams achieve encryption at rest?

level: juniorimportance: must knowfreq 70%
basics
~20 s

No. Open-source Kafka has no built-in encryption of log segment files on disk. Encryption at rest is provided externally: encrypt the broker's volume/disk (LUKS or a cloud KMS-backed encrypted EBS volume), or encrypt the message payload before producing it.

open as a page

What are Kafka's four security protocols, and what does each one provide in terms of encryption and authentication?

level: juniorimportance: must knowfreq 70%
basics
~10 s

Kafka has four security protocols: PLAINTEXT (no encryption, no auth), SSL (TLS encryption, optional client cert auth), SASL_PLAINTEXT (SASL auth, no encryption), and SASL_SSL (SASL auth plus TLS encryption).

open as a page

How do you create and inspect a Kafka topic from the command line, and what do the key options control?

level: juniorimportance: must knowfreq 80%
basics
~10 s

Use kafka-topics.sh with --bootstrap-server. To create: --create --topic name --partitions N --replication-factor R. To inspect: --describe --topic name, which shows partitions, leaders, replicas, and in-sync replicas (ISR).

open as a page

What does Kafka's cleanup.policy control, and how do delete and compact differ? How do retention.ms and retention.bytes fit in?

level: juniorimportance: must knowfreq 70%
basics
~10 s

cleanup.policy sets how old data is removed. delete drops whole log segments once they exceed retention.ms (age) or retention.bytes (size). compact keeps the latest value per key forever. retention.ms/bytes only apply to delete.

open as a page

You start a brand-new Kafka broker with a fresh broker.id and join it to the cluster. Why does it not start serving traffic, and what must you do to put data on it?

level: juniorimportance: must knowfreq 70%
basics
~10 s

A new broker joins the cluster but Kafka never auto-moves existing partitions onto it. You must explicitly reassign replicas to the new broker using kafka-reassign-partitions to scale out.

open as a page

How do you estimate the raw disk storage a Kafka topic (or cluster) will consume given its throughput, retention, and replication settings?

level: juniorimportance: must knowfreq 75%
basics
~20 s

Storage = ingest rate (bytes/sec) x retention seconds x replication factor. So a 10 MB/s topic kept for 7 days at RF=3 needs roughly 10MB x 604800s x 3 = about 18 TB of disk across the cluster.

open as a page

How do you use kafka-configs.sh to alter a topic-level config, and what does the command look like?

level: juniorimportance: must knowfreq 65%
basics
~10 s

Use kafka-configs.sh with --bootstrap-server, --entity-type topics, --entity-name <topic>, and --alter --add-config key=value. To remove an override use --delete-config key. --describe shows current overrides.

open as a page

What is kafka-producer-perf-test.sh, and what do its --num-records, --record-size, and --throughput flags control?

level: juniorimportance: must knowfreq 62%
basics
~10 s

It is Kafka's built-in load generator for producers. --num-records sets how many messages to send, --record-size sets each message's size in bytes, and --throughput caps messages per second (-1 means unthrottled, full speed).

open as a page

What does the UnderReplicatedPartitions broker metric mean, and what should it normally read?

level: juniorimportance: must knowfreq 70%
basics
~20 s

UnderReplicatedPartitions counts how many partitions led by this broker have fewer in-sync replicas than the configured replication factor. In a healthy cluster it should be 0. A sustained non-zero value means replicas are falling behind or down.

open as a page

What is JMX in the context of a Kafka broker, and how do you expose and read broker metrics through it?

level: juniorimportance: must knowfreq 55%
basics
~20 s

JMX (Java Management Extensions) is the Java standard for exposing runtime metrics as MBeans. A Kafka broker publishes its internal metrics as JMX MBeans. You enable it by setting JMX_PORT (or jmxremote system properties) and read it with tools like JConsole, jmxterm, or a Prometheus JMX exporter.

open as a page

What do num.network.threads and num.io.threads control on a Kafka broker, and how do you decide how many to set?

level: juniorimportance: must knowfreq 70%
basics
~10 s

num.network.threads handle reading requests off and writing responses onto network sockets; num.io.threads do the actual work (reading/writing the log on disk). Default 3 and 8. Raise them if request-handler/network idle time drops.

open as a page

What are the most important built-in JMX metrics for a Kafka producer and consumer, and what does each tell you about client health?

level: juniorimportance: must knowfreq 70%
basics
~10 s

Kafka clients expose metrics over JMX. Key producer metrics: record-send-rate, request-latency-avg, buffer-available-bytes. Key consumer metrics: records-consumed-rate, fetch-latency-avg, records-lag-max. They show throughput, latency, and buffering/lag health.

open as a page

In an active-active bidirectional MirrorMaker 2 setup between two clusters, how does the remote-topic naming convention prevent replication loops, and what would happen if you turned it off?

level: juniorimportance: must knowfreq 70%
basics
~20 s

MM2 prefixes replicated topics with the source cluster's alias (e.g. topic 'orders' from cluster A becomes 'A.orders' on cluster B). Because the prefix marks where data came from, MM2 won't replicate a topic back to the cluster it originated from, so records don't loop forever.

open as a page

What is an active-passive (primary/standby) disaster recovery topology in Kafka, and how does it differ from active-active?

level: juniorimportance: must knowfreq 70%
basics
~20 s

In active-passive DR, one Kafka cluster (primary) serves all traffic while a second cluster (standby) only receives a one-way replicated copy. Clients use the standby only after failover. Active-active runs producers/consumers on both clusters at once.

open as a page

What is Confluent Cluster Linking, and how does it differ at a high level from a Connect-based replicator?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Cluster Linking is a broker-built-in feature that copies topics from a source Kafka cluster to a destination cluster over a direct link, with no separate Connect cluster. The destination gets read-only mirror topics that keep the same message offsets.

open as a page

What is 'fetch from follower' in Kafka (KIP-392), and what problem does it solve in a multi-AZ or multi-region deployment?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Normally Kafka consumers read only from the partition leader. Fetch-from-follower (KIP-392) lets a consumer read from a nearby replica (follower) in its own availability zone instead, cutting cross-zone network traffic and the cloud egress cost that comes with it.

open as a page

Is the Kafka KafkaConsumer thread-safe, and what is the canonical threading model for consuming from Kafka?

level: juniorimportance: must knowfreq 70%
basics
~10 s

No. The Java KafkaConsumer is not thread-safe — only one thread may call it (except wakeup()). The standard model is one consumer instance per thread, each in its own poll loop.

open as a page

What is DefaultErrorHandler in Spring for Apache Kafka, and what does it do when a listener throws an exception?

level: juniorimportance: must knowfreq 70%
basics
~20 s

DefaultErrorHandler is the standard component that catches exceptions thrown by your @KafkaListener. It retries the record a configurable number of times with a backoff delay, and if all retries fail it calls a recoverer (by default just logs the failed record).

open as a page

What is the difference between an event and a command in an event-driven system, and how does that distinction shape how you name and publish messages to Kafka topics?

level: juniorimportance: must knowfreq 78%
basics
~20 s

An event states something that already happened (past tense, e.g. OrderPlaced) and has no specific intended recipient. A command tells one service to do something (imperative, e.g. PlaceOrder) and expects an action. Events go to topics for any consumer; commands target one handler.

open as a page

Why does a Kafka consumer often need an application-level deduplication store, and what is the simplest way to build one?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Kafka redelivers messages after rebalances or restarts (at-least-once), so the same record can arrive more than once. You record each processed message's id in a store (a DB unique constraint or Redis key) and skip any id you've already seen.

open as a page

How do you construct a KafkaProducer and a KafkaConsumer in the Java client, and what are the minimum required configuration properties for each?

level: juniorimportance: must knowfreq 75%
basics
~20 s

Build a Properties (or Map) of config, then pass it to the constructor: new KafkaProducer<>(props) or new KafkaConsumer<>(props). Producers need bootstrap.servers plus a key and value serializer; consumers need bootstrap.servers plus a key and value deserializer (and usually group.id).

open as a page

What is the fundamental architectural difference between how Apache Pulsar and Apache Kafka store and serve data?

level: juniorimportance: must knowfreq 70%
basics
~10 s

Kafka couples storage and serving: each broker owns its partition data on local disk. Pulsar separates them: brokers are stateless serving nodes, and storage lives in a separate layer (Apache BookKeeper).

open as a page

What is librdkafka, and what is its relationship to the Kafka clients available for Python, Go, .NET, and Node.js?

level: juniorimportance: must knowfreq 70%
basics
~20 s

librdkafka is a C/C++ implementation of the Kafka client protocol. The official Python, Go, .NET, and Node.js clients are thin language bindings that wrap librdkafka, so they share one battle-tested core instead of each reimplementing Kafka.

open as a page

What is Kafka Connect, and how do sink connectors get data from Kafka topics into a data lake or warehouse like S3, BigQuery, or Snowflake?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Kafka Connect is a framework for moving data between Kafka and external systems without custom code. A sink connector reads records from Kafka topics and writes them to a destination like S3, BigQuery, or Snowflake using ready-made plugins and config.

open as a page

What is the Azure Event Hubs Kafka endpoint, and how does an existing Kafka application connect to it?

level: juniorimportance: must knowfreq 55%
basics
~20 s

Azure Event Hubs exposes a Kafka-compatible endpoint on port 9093, so a normal Kafka client can produce and consume by just changing bootstrap.servers and using SASL/SSL auth with a connection string — no Kafka broker is actually run.

open as a page

What are Redpanda and WarpStream, and what does it mean that they are 'Kafka-compatible'?

level: juniorimportance: must knowfreq 55%
basics
~20 s

They are alternative streaming systems that speak the Kafka wire protocol, so existing Kafka clients work unchanged. Redpanda is a C++ broker with no JVM or ZooKeeper; WarpStream stores data in S3 with stateless brokers.

open as a page
Apache Kafka interview questions & primer · KataJob