Apache Kafka
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 pageshowhide
guide
overview
~1 minApache 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.
- Topics, Partitions and Log Storage →
The data model every other section assumes: partitions, keys, ordering, offsets and retention.
- Producers →
The write path, where acks, batching and idempotence decide whether data is lost or duplicated.
- Consumers and Consumer Groups →
The read path: groups, rebalancing and offset commits, where most production incidents first show up.
- Replication and Durability →
What an acknowledged write really survives once brokers fail, and the durability-versus-availability trade.
- Delivery Semantics and Transactions →
Delivery guarantees and transactions make sense only after producers, consumers and replication are clear.
- 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
- Architecture and Internals45 questions
- 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
- Cluster Coordination: KRaft and ZooKeeper61 questions
- ZooKeeper's Historical Role5 questions
- ZooKeeper Scalability Limits5 questions
- KRaft Architecture and Metadata Quorum5 questions
- Raft Protocol and Leader Election5 questions
- Active Controller and Metadata Log5 questions
- Metadata Snapshots and Log Replay5 questions
- Broker Registration and Heartbeats5 questions
- Dynamic Quorum Reconfiguration5 questions
- ZooKeeper-to-KRaft Migration6 questions
- KRaft Storage Format and Cluster Bootstrap5 questions
- KRaft Operations and Observability5 questions
- Topics, Partitions and Log Storage47 questions
- Topics and Partition Fundamentals5 questions
- Partition Keys and Hash Routing5 questions
- Ordering Guarantees5 questions
- Choosing and Changing Partition Count5 questions
- Retention Policies and Cleanup (Delete)5 questions
- Log Compaction and Tombstones5 questions
- Offset Addressing Model6 questions
- Record Batches and Message Format6 questions
- Topic Configuration and Administration5 questions
- Producers36 questions
- Send Flow and Record Accumulator5 questions
- Acks and Durability5 questions
- Batching, Linger and Compression5 questions
- Idempotent Producer5 questions
- Partitioner Strategies5 questions
- Error Handling, Retries and Ordering5 questions
- Serializers and Interceptors6 questions
- Consumers and Consumer Groups55 questions
- Poll Loop and Fetch Mechanics5 questions
- Heartbeat and Session Liveness5 questions
- Rebalancing and Assignment Strategies5 questions
- Offset Management and Commit Strategies5 questions
- Position, Seek and Offset Reset5 questions
- Consumer Lag and Measurement5 questions
- Static Membership5 questions
- Deserialization and Consumer Configuration5 questions
- Replication and Durability50 questions
- In-Sync Replicas and Replica Lag5 questions
- Unclean Leader Election5 questions
- Rack Awareness and Replica Placement5 questions
- End-to-End Durability vs Availability Tuning5 questions
- Leadership Failover and Replica State5 questions
- Delivery Semantics and Transactions40 questions
- At-Most/At-Least/Exactly-Once Defined5 questions
- Transactional Producer API5 questions
- Transaction Coordinator and Zombie Fencing5 questions
- Transaction Markers and Atomic Commit5 questions
- Read-Committed Isolation and LSO5 questions
- Consume-Transform-Produce EOS5 questions
- Idempotent Consumers and Dedup5 questions
- EOS Scope and Limits5 questions
- Kafka Streams61 questions
- Streams DSL and Topology5 questions
- Stateless Operations6 questions
- State Stores and Fault Tolerance6 questions
- Joins5 questions
- Windowing5 questions
- Time Semantics5 questions
- Aggregations6 questions
- Processor API and Punctuators6 questions
- Interactive Queries5 questions
- Exactly-Once Processing6 questions
- Scaling, Threads and Tasks6 questions
- Kafka Connect55 questions
- Connect Framework and Runtime6 questions
- Connect REST API and Lifecycle Management6 questions
- Source vs Sink Connectors and Tasks5 questions
- Converters and Schema Handling6 questions
- Single Message Transforms and Predicates5 questions
- Offset Management in Connect5 questions
- Error Handling and Dead-Letter Queue5 questions
- Change Data Capture with Debezium6 questions
- Scaling and Rebalancing Across Workers6 questions
- Connect Security and Config Externalization5 questions
- Schemas and Serialization48 questions
- Confluent Schema Registry6 questions
- Avro With Kafka5 questions
- Protobuf and JSON Schema Support5 questions
- Schema Evolution and Compatibility Modes5 questions
- Subject Naming Strategies5 questions
- Schema References, Contexts and Governance6 questions
- Security62 questions
- TLS/SSL Encryption in Transit5 questions
- Mutual TLS Client Authentication5 questions
- SASL PLAIN and SCRAM Authentication5 questions
- SASL Kerberos (GSSAPI)5 questions
- SASL OAUTHBEARER and Delegation Tokens5 questions
- Authorization and ACLs6 questions
- Super Users and Authorizer Defaults5 questions
- Security Protocol and Listener Matrix6 questions
- Quotas as a Control Plane5 questions
- Encryption at Rest and Governance5 questions
- ZooKeeper and KRaft Metadata Hardening5 questions
- Security Auditing and Monitoring5 questions
- Operations and Administration46 questions
- Cluster Sizing and Capacity Planning5 questions
- Broker and Topic Configuration Management6 questions
- Partition Reassignment and Data Balancing5 questions
- Adding, Removing and Decommissioning Brokers5 questions
- Rolling Upgrades and Version Compatibility5 questions
- Tiered and Remote Storage5 questions
- Backup, Disaster Recovery and Retention Ops5 questions
- Admin CLI and AdminClient5 questions
- Operational Monitoring and Alerting5 questions
- Monitoring and Performance Tuning54 questions
- Broker JMX Metrics and Health Signals6 questions
- Producer Throughput vs Latency Tuning6 questions
- Consumer Fetch and Parallelism Tuning5 questions
- Broker Threads, Page Cache and Disk Tuning5 questions
- Benchmarking with Perf Test Tools6 questions
- Client Metrics Instrumentation and Reporters5 questions
- Multi-Cluster and Geo-Replication65 questions
- MirrorMaker 2 Replication Flows5 questions
- MM2 Offset Translation and Checkpoints6 questions
- MM2 Heartbeats and Replication Monitoring5 questions
- Remote Topic Naming and Replication Policies6 questions
- Cluster Linking and Mirror Topics5 questions
- Active-Passive Disaster Recovery Topologies6 questions
- Active-Active Bidirectional Replication5 questions
- Stretch Clusters and Rack-Aware Placement5 questions
- Fetch From Follower and Cross-Region Reads5 questions
- Replication Topology Pitfalls and Cycles5 questions
- Multi-Region Design Tradeoffs6 questions
- Cross-Cluster Security and Connectivity6 questions
- Client Development and Integration Patterns62 questions
- Java Client Configuration and Lifecycle6 questions
- Error Handling, Retry Topics and DLT6 questions
- Transactional Outbox Pattern6 questions
- Idempotent Consumers and Deduplication5 questions
- Event-Driven Design with Kafka6 questions
- Testing Kafka Applications6 questions
- Consumer Concurrency and Threading Models5 questions
- Reactive and Async Kafka Clients6 questions
- Request-Reply and Correlation Patterns5 questions
- Ecosystem and Platform Choices48 questions
- ksqlDB Streaming SQL6 questions
- Client and API Ecosystem5 questions
- Managed Kafka Offerings5 questions
- Kafka-Compatible Alternatives5 questions
- Apache Pulsar Comparison5 questions
- Kafka vs Traditional Message Brokers5 questions
- Stream Processing Engine Choice5 questions
- Kafka in the Data Ecosystem6 questions
- Cloud Event-Hubs and Protocol-Edge Choices6 questions
questions
835 · 16 sectionsWhat is the purpose of bootstrap.servers, and why don't you need to list every broker?
basics
~20 sbootstrap.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.
What is a Kafka broker, and what role does broker.id play in a cluster?
basics
~10 sA 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.
What is a Kafka partition's commit log, and why is it split into segments on disk?
basics
~20 sEach 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.
What is the Kafka log cleaner, and how do you enable it for a topic?
basics
~20 sThe 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.
Why does Kafka rely on the operating system page cache instead of maintaining its own in-process (JVM heap) record cache?
basics
~20 sKafka 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.
Cluster Coordination: KRaft and ZooKeeper
all 61 Cluster Coordination: KRaft and ZooKeeper questions →In KRaft mode, what is the active controller and what role does it play in managing cluster metadata?
basics
~20 sThe 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.
What is KRaft in Apache Kafka, and what problem was it introduced to solve?
basics
~10 sKRaft (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.
In a KRaft Kafka cluster, how does a broker tell the controller it is still alive, and what RPC is involved?
basics
~10 sEach 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.
What is the difference between the static controller.quorum.voters config and the dynamic quorum introduced by KIP-853?
basics
~10 sStatic 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.
In a KRaft cluster, why does Kafka periodically take metadata snapshots of the __cluster_metadata log?
basics
~20 sThe 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.
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?
basics
~10 sThe 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.
What does cleanup.policy=compact do to a Kafka topic, and how is it different from the default delete policy?
basics
~10 sCompaction 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.
What is an offset in Kafka, and why are offsets meaningful only within a single partition?
basics
~20 sAn 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.
What ordering guarantee does Kafka provide, and what does it NOT guarantee?
basics
~10 sKafka 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.
What does `kafka-topics.sh --alter --partitions` allow you to do, and what is the key directionality constraint?
basics
~10 sIt 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.
What do the producer acks settings 0, 1, and all (-1) mean, and how do they trade durability against latency?
basics
~20 sacks=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.
What do batch.size and linger.ms control in a Kafka producer, and how do they trade latency for throughput?
basics
~20 sbatch.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.
What is the difference between a retriable and a fatal (non-retriable) exception in the Kafka producer, and how does each affect a send?
basics
~20 sRetriable 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.
What is Kafka's idempotent producer, and how do you turn it on?
basics
~10 sAn 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).
How does a Kafka producer decide which partition a record goes to when the record has a non-null key?
basics
~20 sFor 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.
What are key.deserializer and value.deserializer in a Kafka consumer, and why must they pair with the producer's serializers?
basics
~10 sKafka 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.
What is the Group Coordinator in Kafka, and how is the coordinator broker for a particular consumer group chosen?
basics
~10 sThe 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.
What is a consumer heartbeat in Kafka, and what does it accomplish?
basics
~20 sA 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.
How do you inspect consumer lag from the command line, and what do the columns of kafka-consumer-groups.sh --describe mean?
basics
~10 sRun 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.
What is consumer lag in Kafka, and how is it calculated for a single partition?
basics
~10 sConsumer 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.
What are the three main Kafka settings that work together to control write durability, and what does each do at a high level?
basics
~20 sReplication 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.
When the broker hosting a partition's leader replica dies, what happens to that partition so producers and consumers can keep working?
basics
~10 sKafka 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.
After a Kafka write returns successfully, where does the data physically live, and why does that distinction matter?
basics
~20 sAfter 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.
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?
basics
~20 sLEO 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.
What is the In-Sync Replicas (ISR) set in Apache Kafka, and how does it relate to the assigned replicas of a partition?
basics
~20 sThe 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.
What is the consume-transform-produce pattern in Kafka, and why does plain at-least-once delivery fall short of exactly-once for it?
basics
~20 sIt'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.
What is the scope boundary of Kafka's exactly-once semantics (EOS), and why doesn't it automatically extend to external systems?
basics
~20 sKafka 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.
What are the three delivery semantics in Kafka (at-most-once, at-least-once, exactly-once), and what does each guarantee about message delivery?
basics
~20 sAt-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).
What does it mean to make a Kafka consumer idempotent, and why does that let you live safely with at-least-once delivery?
basics
~20 sAn 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.
What does the consumer setting isolation.level do in Kafka, and what are its two possible values?
basics
~10 sisolation.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.
In Kafka Streams, what is the difference between groupByKey() and groupBy(), and why does one of them trigger a repartition?
basics
~20 sgroupByKey() 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.
What is the difference between a KStream and a KTable in the Kafka Streams DSL, and when would you use each?
basics
~20 sA 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.
What does setting processing.guarantee=exactly_once_v2 in a Kafka Streams application actually guarantee, and how do you enable it?
basics
~20 sIt 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.
What are Interactive Queries in Kafka Streams, and how do you read the value for a key from a local state store?
basics
~20 sInteractive 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().
What kinds of joins does Kafka Streams support, and how do they differ at a high level?
basics
~20 sKafka 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.
What is Change Data Capture (CDC) with Debezium, and why is log-based CDC preferred over query-based polling?
basics
~20 sDebezium 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.
What is a converter in Kafka Connect, and what is the difference between key.converter and value.converter?
basics
~10 sA 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.
What does errors.tolerance control in Kafka Connect, and what is the difference between the values none and all?
basics
~10 serrors.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.
What is Kafka Connect, and what problem does it solve compared to writing your own producer/consumer applications?
basics
~10 sKafka 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.
How does Kafka Connect track offsets for source connectors versus sink connectors? Where is each kind stored?
basics
~20 sSource 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.
What is Apache Avro and why is it commonly used as the serialization format for Kafka records?
basics
~20 sAvro 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.
What are the compatibility modes in Confluent Schema Registry, and at a high level what does each one allow?
basics
~20 sSchema 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.
Why is deserializing a Kafka message treated as security-sensitive, and what is the core threat when a consumer deserializes an untrusted payload?
basics
~20 sDeserialization 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.
What is a "poison pill" record in a Kafka consumer, and why can it stall a consumer group?
basics
~20 sA 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.
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?
basics
~10 sUse 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.
How do you turn on an audit log of allowed and denied authorization decisions in Apache Kafka, and where does that output go?
basics
~20 sKafka'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.
What is a Kafka ACL, and what are the fields that make up a single ACL binding?
basics
~20 sAn 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.
What is the super.users setting in Kafka, and what happens when a principal is listed in it?
basics
~10 ssuper.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.
Does Apache Kafka natively encrypt the data it writes to disk (log segments)? If not, how do teams achieve encryption at rest?
basics
~20 sNo. 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.
What are Kafka's four security protocols, and what does each one provide in terms of encryption and authentication?
basics
~10 sKafka 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).
How do you create and inspect a Kafka topic from the command line, and what do the key options control?
basics
~10 sUse 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).
What does Kafka's cleanup.policy control, and how do delete and compact differ? How do retention.ms and retention.bytes fit in?
basics
~10 scleanup.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.
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?
basics
~10 sA 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.
How do you estimate the raw disk storage a Kafka topic (or cluster) will consume given its throughput, retention, and replication settings?
basics
~20 sStorage = 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.
How do you use kafka-configs.sh to alter a topic-level config, and what does the command look like?
basics
~10 sUse 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.
What is kafka-producer-perf-test.sh, and what do its --num-records, --record-size, and --throughput flags control?
basics
~10 sIt 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).
What does the UnderReplicatedPartitions broker metric mean, and what should it normally read?
basics
~20 sUnderReplicatedPartitions 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.
What is JMX in the context of a Kafka broker, and how do you expose and read broker metrics through it?
basics
~20 sJMX (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.
What do num.network.threads and num.io.threads control on a Kafka broker, and how do you decide how many to set?
basics
~10 snum.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.
What are the most important built-in JMX metrics for a Kafka producer and consumer, and what does each tell you about client health?
basics
~10 sKafka 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.
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?
basics
~20 sMM2 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.
What is an active-passive (primary/standby) disaster recovery topology in Kafka, and how does it differ from active-active?
basics
~20 sIn 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.
What is Confluent Cluster Linking, and how does it differ at a high level from a Connect-based replicator?
basics
~20 sCluster 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.
When you replicate data between two Kafka clusters with MirrorMaker 2, what does it mean to 'secure the replication link', and what are the two layers involved?
basics
~20 sMirrorMaker 2 is just a Kafka client. Securing the link means: (1) encrypt traffic in transit with TLS, and (2) authenticate MM2 to each cluster (mTLS or SASL) so only an authorized identity can read the source and write the target.
What is 'fetch from follower' in Kafka (KIP-392), and what problem does it solve in a multi-AZ or multi-region deployment?
basics
~20 sNormally 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.
Client Development and Integration Patterns
all 62 Client Development and Integration Patterns questions →Is the Kafka KafkaConsumer thread-safe, and what is the canonical threading model for consuming from Kafka?
basics
~10 sNo. 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.
What is DefaultErrorHandler in Spring for Apache Kafka, and what does it do when a listener throws an exception?
basics
~20 sDefaultErrorHandler 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).
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?
basics
~20 sAn 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.
Why does a Kafka consumer often need an application-level deduplication store, and what is the simplest way to build one?
basics
~20 sKafka 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.
How do you construct a KafkaProducer and a KafkaConsumer in the Java client, and what are the minimum required configuration properties for each?
basics
~20 sBuild 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).
What is the fundamental architectural difference between how Apache Pulsar and Apache Kafka store and serve data?
basics
~10 sKafka 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).
What is librdkafka, and what is its relationship to the Kafka clients available for Python, Go, .NET, and Node.js?
basics
~20 slibrdkafka 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.
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?
basics
~20 sKafka 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.
What is the Azure Event Hubs Kafka endpoint, and how does an existing Kafka application connect to it?
basics
~20 sAzure 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.
What are Redpanda and WarpStream, and what does it mean that they are 'Kafka-compatible'?
basics
~20 sThey 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.