Kafka Consumer Groups and Rebalancing Explained
How Kafka consumer groups share partitions, what triggers a rebalance, and how cooperative assignment and static membership cut downtime.
Event streaming interviews: topics and partitions, replication, producers and acks, consumer groups and rebalancing, offsets, delivery guarantees and exactly-once.
Official reference: Apache Kafka documentation
How Kafka consumer groups share partitions, what triggers a rebalance, and how cooperative assignment and static membership cut downtime.
How Kafka transactions make consume-transform-produce exactly-once, how zombie fencing and read_committed work, and where the guarantee ends.
How the Kafka idempotent producer stops duplicates from retries, what it cannot dedupe, and how batch.size, linger.ms and compression work.
How Kafka consumers commit offsets, why commit order decides at-most-once or at-least-once, and how to get exactly-once effects.
How Kafka partitions and offsets work, why ordering is only per partition, and how keys keep each entity's events in sequence.
Why one bad Kafka record can block a whole partition, and how to handle it with error classification, bounded retries and a dead letter topic.
How Kafka replicates partitions, what the ISR and high watermark mean, and how acks=all with min.insync.replicas prevents data loss.
How Kafka deletes old segments by time or size, how log compaction keeps the latest value per key, and how tombstones delete keys.
Kafka is a distributed, replicated commit log. Producers append records to topics, the brokers store them durably for a configured retention period, and consumers read them by position (offset).
The differences from a classic queue:
That makes Kafka a good fit for event streaming, high-throughput pipelines and event sourcing, while a classic broker is simpler for per-message routing, priorities and task queues. Kafka 4.0 adds early-access share groups (queue-like consumption), but the log model is still the core.
Likely follow-up: When would you still choose RabbitMQ or SQS over Kafka? · What does it cost to keep data for 30 days instead of 7?
A topic is a named stream of records, such as orders. It is split into partitions, and each partition is an ordered, append-only log stored on a broker (with replicas on other brokers).
An offset is the position of a record inside one partition: 0, 1, 2 and so on. Offsets are only meaningful per partition, so orders-0 offset 42 and orders-3 offset 42 are unrelated records.
Partitions are the unit of parallelism and ordering. Producers choose a partition (usually by hashing the key), consumers in a group each own a subset of partitions, and Kafka guarantees order only inside a partition. A consumer's progress is stored as a committed offset per partition: the offset of the next record it should read, not the last one it processed.
Likely follow-up: Can two records in different partitions have the same offset? · Do offsets ever get reused after retention deletes data?
Kafka guarantees order within a partition only. Records in one partition are stored and delivered in offset order; there is no ordering across partitions of a topic.
To keep all events for one entity in order, give them the same key, for example the order ID. The default partitioner hashes the key (murmur2) modulo the partition count, so every event for order-1 lands in the same partition and is read by one consumer in sequence.
Three things can still break it:
Likely follow-up: How would you get a global order across a whole topic, and what does it cost?
A consumer group is a set of consumers sharing a group.id that cooperate to read a topic. Each partition is assigned to exactly one member of the group at a time, so the group processes every record once between them, and the group's committed offsets are stored in the internal __consumer_offsets topic.
Different groups are independent: a billing group and an analytics group each receive every record and keep their own offsets.
If there are more consumers than partitions, the extra consumers get no partitions and sit idle. They are not useless, because they take over when a member fails, but they add no throughput. The maximum parallelism of one group on one topic is the partition count, which is why partition count is a capacity decision.
Likely follow-up: How do you scale processing beyond the partition count?
A cluster is a set of brokers that store partition replicas and serve producers and consumers, plus a controller that manages metadata: which topics exist, which broker leads each partition, which replicas are in sync, and broker membership.
Historically the controller kept that metadata in ZooKeeper, a separate system to deploy and secure. KRaft (Kafka Raft) replaces it: a small quorum of controller nodes, usually three, replicates metadata in an internal Raft log called __cluster_metadata. Brokers follow that log instead of being pushed updates.
Timeline: KRaft became production-ready in 3.3, ZooKeeper mode was deprecated in 3.5, 3.9 is the last release that supports it (the bridge for migration), and Kafka 4.0 runs only in KRaft mode. Benefits are one system to operate, faster controller failover and much higher partition limits.
Likely follow-up: What is the difference between combined and dedicated controller mode?
Each partition is stored as a series of segment files. New records go to the active segment; when it reaches segment.bytes (1 GiB by default) or segment.ms (7 days) it is closed and a new one starts.
With cleanup.policy=delete (the default), a closed segment is deleted once its newest record is older than retention.ms (7 days by default) or when the partition exceeds retention.bytes (unlimited by default). Deletion is per segment, never per record, and the active segment is never deleted, so data can live somewhat longer than the retention setting.
With cleanup.policy=compact, Kafka keeps the latest record per key instead, which suits changelog topics. Reading has no effect on retention: a record is kept for the full period whether ten groups read it or none.
Likely follow-up: What happens to a consumer whose committed offset has been deleted by retention?
acks=0, acks=1 and acks=all mean?easyacks controls when the partition leader answers a produce request, which is the trade-off between latency and durability:
acks=0: the producer does not wait for any response. Fastest, but records can be lost silently and retries cannot work.acks=1: the leader replies after writing to its own log. If the leader dies before followers copy the record, it is lost.acks=all (or -1): the leader replies only after every replica currently in the in-sync replica set has the record. This is the default since Kafka 3.0.acks=all alone is not enough if the ISR has shrunk to just the leader, so it is paired with the topic setting min.insync.replicas, typically 2 with replication factor 3. Then a write is accepted only when at least two replicas have it.
Likely follow-up: Does acks=all make a write slower if one follower is slow but still in sync?
auto.offset.reset do, and when does it apply?easyauto.offset.reset tells a consumer where to start when there is no valid committed offset for a partition. It does not apply on every restart.
It is used in two cases: a brand-new group that has never committed, and a committed offset that is out of range, typically because retention deleted that data while the group was down.
The values:
latest (default): start at the end and read only new records.earliest: start at the oldest retained record.none: throw an exception so the application decides.A classic surprise is deploying a new service with the default latest and missing every event produced before its first poll. Another is a consumer that falls behind retention and silently jumps, either skipping data with latest or reprocessing everything with earliest. To replay deliberately, reset the group's offsets with kafka-consumer-groups.sh --reset-offsets while the group is inactive.
Likely follow-up: How would you replay the last two hours of a topic for one group?
Each partition has a replication factor, typically 3. One replica is the leader: all produce requests, and by default all fetches, go to it. The followers fetch from the leader continuously, like consumers.
The in-sync replica set (ISR) is the leader plus followers that have caught up recently. A follower that has not caught up within replica.lag.time.max.ms (30 seconds by default) is removed from the ISR, and rejoins once it catches up.
The high watermark is the highest offset copied to every ISR member. Records below it are committed; consumers only see records below the high watermark, so they never read data that could vanish in a failover.
When a leader fails, the controller elects a new leader from the ISR, so no committed record is lost. Leader epochs let a returning replica truncate any uncommitted tail that diverged from the new leader.
Likely follow-up: What causes ISR shrink events, and how would you alert on them?
acks=all and min.insync.replicas=2, what happens as brokers fail one by one?midmin.insync.replicas is the minimum ISR size for which a broker accepts an acks=all write.
acks=all producers get NotEnoughReplicasException and retry until delivery.timeout.ms. The partition is effectively read-only for them; consumers can still read.This is the deliberate trade: availability for durability. Producers using acks=1 would still succeed, because min.insync.replicas only applies to acks=all. RF 3 with a minimum of 2 is the common choice because it tolerates one failure, or one rolling restart, without blocking writes.
Likely follow-up: Why is min.insync.replicas equal to the replication factor a bad idea?
Normally a new leader must come from the ISR. If every in-sync replica is unavailable, the partition goes offline until one of them returns.
unclean.leader.election.enable=true allows an out-of-sync replica to become leader instead. The partition comes back sooner, but any records the old leader had that this replica never copied are lost, and when the old leader returns it truncates its log to match. Consumers may also see offsets reused for different data.
The default has been false since Kafka 0.11, and it should stay off for anything where losing acknowledged data is unacceptable: payments, orders, event-sourced state. You might enable it per topic for metrics or logs where availability matters more than completeness. It can also be triggered once, manually, with kafka-leader-election.sh --election-type unclean during an outage, which is safer than leaving it on.
Likely follow-up: How do rack-aware replica placement and min.insync.replicas reduce the need for it?
Without idempotence, a retry can create a duplicate: the broker writes a batch, the acknowledgement is lost, and the producer sends it again. With several requests in flight, a retry can also reorder records.
With enable.idempotence=true (the default since Kafka 3.0), the broker assigns the producer a producer ID, and the producer numbers every batch with a sequence number per partition. The broker remembers the last sequences per producer and partition: a batch it has already written is acknowledged without being appended again, and an out-of-order batch is rejected so the producer can resend in order. This works with up to 5 in-flight requests per connection.
The limits matter: it deduplicates retries inside one producer session on one partition. If the application crashes and resends, or two instances send the same event, the new producer ID makes those different records. Cross-session guarantees need transactions or consumer-side deduplication.
Likely follow-up: What happens to idempotence if you explicitly set acks=1?
batch.size, linger.ms and compression.type control?midsend() does not make a network call. It serializes the record, picks a partition and appends it to an in-memory batch for that partition in the record accumulator. A background sender thread ships batches to partition leaders.
batch.size (16 KB default) is the maximum bytes per partition batch; a full batch is sent immediately.linger.ms is how long the producer waits for a batch to fill before sending it anyway. The default changed from 0 to 5 ms in Kafka 4.0; raising it to 10 to 50 ms often multiplies throughput for a small latency cost.compression.type (none, gzip, snappy, lz4, zstd) compresses whole batches, so bigger batches compress better.buffer.memory (32 MB) caps unsent data; when it is full, send() blocks up to max.block.ms.Small batches mean many requests and poor compression. The usual tuning is bigger batches with some linger and lz4 or zstd.
Likely follow-up: Why can calling .get() on every send() destroy throughput?
If the record specifies a partition, that is used. Otherwise:
A hot partition appears when keys are skewed: one huge tenant, a default key such as "unknown", or a low-cardinality key like a country code. One partition then carries most of the traffic and one consumer lags while others idle. Fixes include a better key, salting a hot key into several sub-keys when per-key order allows it, or splitting the heavy tenant into its own topic.
Likely follow-up: Why is a custom partitioner rarely the right first fix?
A rebalance redistributes partitions among group members. It is triggered when a member joins, leaves or is considered dead (missed heartbeats for session.timeout.ms, 45 seconds by default, or did not call poll() within max.poll.interval.ms, 5 minutes), or when partitions are added to a subscribed topic.
With the classic protocol there are two styles:
CooperativeStickyAssignor): members keep the partitions they still own and give up only those that must move, over two shorter rounds. Most consumers never stop.Kafka 4.0 also makes the new protocol from KIP-848 generally available (group.protocol=consumer): the broker computes assignments and reconciles each member incrementally, removing the group-wide barrier entirely.
Likely follow-up: How do you migrate a running group from eager to cooperative assignment?
By default every consumer gets a new, dynamic member ID when it joins. A restart therefore looks like one member leaving and another joining, which causes two rebalances per pod during a rolling deploy.
Static membership gives each consumer a stable group.instance.id, such as the StatefulSet pod name. The coordinator remembers that identity and its assignment. If the instance comes back within session.timeout.ms, it gets the same partitions with no rebalance, and a static member does not send a leave request when it shuts down.
The trade-off is slower failure detection: a crashed static member's partitions sit unconsumed until the session timeout expires, so teams raise session.timeout.ms to cover a restart but not much longer. Each instance ID must be unique; two live consumers with the same ID cause one to be fenced with an error. It pairs well with cooperative assignment and with local state such as Kafka Streams stores.
Likely follow-up: What happens if two pods start with the same group.instance.id?
commitSync and commitAsync. When does each lose or duplicate messages?midenable.auto.commit=true, every 5 seconds): the consumer commits the positions from earlier polls during later poll() calls. If you process synchronously in the poll loop, a crash replays up to 5 seconds of records, which is at-least-once. If you hand records to another thread, the commit can run before processing finishes, and a crash loses them.commitSync() after processing a batch: blocks and retries until it succeeds or fails fatally. Clear at-least-once semantics, with some latency per commit.commitAsync(): does not block or retry, because a retried old commit could overwrite a newer one. Use a callback to log failures.The common production pattern is commitAsync() in the loop for speed and one commitSync() on shutdown and in the onPartitionsRevoked callback, so the last progress is saved before partitions move. Committing before processing gives at-most-once.
Likely follow-up: Why does commitAsync deliberately not retry?
The guarantee is decided by where the consumer commits relative to processing and by producer settings.
isolation.level=read_committed. Kafka Streams enables this with processing.guarantee=exactly_once_v2.The scope matters: Kafka transactions do not cover a database write or an email. For external side effects you get exactly-once effects by making them idempotent, such as an upsert keyed by event ID or a processed-events table updated in the same database transaction.
Likely follow-up: How would you get exactly-once effects when the consumer writes to PostgreSQL?
Lag is, per partition, the distance between the latest offset in the log and the group's committed offset. Total lag is the sum across partitions, but the per-partition view matters: one hot partition can lag while the rest are idle.
Measure it with kafka-consumer-groups.sh --describe --group <name> (columns CURRENT-OFFSET, LOG-END-OFFSET, LAG), the consumer's records-lag-max metric, or an exporter feeding Prometheus. Alert on lag growing over time, or on time lag, rather than on a fixed number.
To reduce it, first find where time goes:
max.poll.interval.ms.Likely follow-up: Why can lag show zero while the service is still behind on business work?
With cleanup.policy=compact, Kafka guarantees to keep at least the latest record for each key instead of deleting by age. A background cleaner rewrites closed segments, dropping older records whose key has a newer value. The active segment is never compacted, and offsets never change: compaction leaves gaps.
A tombstone is a record with a key and a null value. It marks the key as deleted. Compaction removes earlier values for that key, keeps the tombstone for delete.retention.ms (24 hours by default) so slow consumers can see the deletion, then removes it too.
Compacted topics are used for changelogs and current state: Kafka Streams state store changelogs, Connect offsets, __consumer_offsets and CDC topics. Compaction is not instant, so consumers may see several values for a key, and records with null keys are rejected. compact,delete combines both policies.
Likely follow-up: How can you bound how long an old value may survive before compaction?
A poison message is a record that fails every time, such as bad JSON, an unknown schema ID or data that violates a business rule. Because the consumer never commits past it, it restarts at the same offset and fails again, so the whole partition is blocked while lag grows.
Handle it by classifying errors:
RecordDeserializationException with the partition and offset, so you can record it and seek() past it. Spring Kafka offers ErrorHandlingDeserializer.Then treat the DLT as an operational queue: alert on it, fix, and replay deliberately.
Likely follow-up: How do retry topics affect per-key ordering?
Kafka stores bytes and does not validate them. A schema registry, such as Confluent Schema Registry, stores versioned Avro, Protobuf or JSON Schema definitions per subject (by default <topic>-value and <topic>-key). The serializer registers or looks up the schema and prefixes each message with a magic byte and a 4-byte schema ID, and the deserializer fetches that schema to decode it.
The registry checks every new version against a compatibility rule:
_TRANSITIVE variants check against all earlier versions, not just the latest.Renaming a field or changing its type usually breaks compatibility.
Likely follow-up: Why can BACKWARD (non-transitive) still break a consumer reading old data from a long-retention topic?
Kafka Connect is a framework for moving data between Kafka and other systems using configuration instead of code. Source connectors pull data in (for example Debezium reading a PostgreSQL write-ahead log), and sink connectors push data out (to S3, Elasticsearch or a JDBC database).
Its pieces:
errors.tolerance=all with errors.deadletterqueue.topic.name to divert bad records.Use Connect when a maintained connector exists for a standard integration. Write your own client when the logic is business processing rather than data movement.
Likely follow-up: How does Debezium capture changes without polling the table?
KStream and a KTable?midKafka Streams is a Java library, not a separate cluster, for building stream processing applications. It runs inside your service, uses consumer groups for scaling (the application.id is the group ID), and parallelizes by input partition into tasks.
KStream is an unbounded sequence of independent events: every record is a fact, such as a payment made.KTable is a changelog view: each record is an update to the current value for its key, like a table row being upserted. A null value deletes the key.Stateful operations (aggregations, joins, windows) keep data in local state stores, RocksDB by default, each backed by a compacted changelog topic so a new instance can rebuild the state. Re-keying operations create internal repartition topics. Setting processing.guarantee=exactly_once_v2 makes processing, state updates and output atomic using transactions. A GlobalKTable replicates a whole topic to every instance for joins without co-partitioning.
Likely follow-up: What is co-partitioning, and why do joins require it?
read_committed change?hardA producer with a transactional.id calls initTransactions() once. The transaction coordinator gives it a producer ID and bumps an epoch, which fences any older instance with the same ID.
Per batch the loop is: beginTransaction(), send the output records, call sendOffsetsToTransaction(offsets, consumer.groupMetadata()) so the input offsets become part of the transaction, then commitTransaction(). On failure it calls abortTransaction() and rewinds the consumer. The coordinator logs state in __transaction_state and writes commit or abort markers into every partition touched.
Records are written before the commit, so visibility is controlled on the read side. Consumers with isolation.level=read_committed read only up to the last stable offset, the point before the earliest open transaction, and skip aborted records. The default read_uncommitted sees everything, including aborted data. Output, offsets and fencing together give exactly-once for Kafka-to-Kafka processing.
Likely follow-up: Why does a long-running open transaction increase end-to-end latency for read_committed consumers?
No. Kafka transactions make writes to Kafka and offset commits atomic. A PostgreSQL insert is outside that transaction, so the classic failure remains: the insert commits, the consumer crashes before its offset commit, and the record is processed again after the rebalance.
The fix is to make the database the source of truth for progress or for deduplication:
INSERT ... ON CONFLICT DO NOTHING, so a replay has no effect.seek() to the stored offset instead of using Kafka's committed offset.Either gives exactly-once effects on top of at-least-once delivery. A two-phase commit across Kafka and the database is not supported by the Java client and is rarely worth it.
Likely follow-up: How would you handle a consumer that calls a third-party payment API?
OrderCreated to Kafka. What can go wrong, and how does the outbox pattern fix it?hardThis is the dual-write problem. The database commit and the Kafka send are two independent operations:
With the transactional outbox, the service writes the order and an outbox row (event ID, aggregate ID, type, payload) in one database transaction. A separate relay publishes outbox rows to Kafka: either a polling publisher, or change data capture such as Debezium reading the write-ahead log, with its outbox event router transform.
The relay is at-least-once: it can publish a row and crash before recording that, so consumers must deduplicate by event ID. Using the aggregate ID as the Kafka key keeps each order's events in sequence. Clean up published outbox rows to keep the table small.
Likely follow-up: Polling publisher or CDC: what are the trade-offs?
Start from throughput and consumer parallelism. Measure what one partition sustains for your producer and for your slowest consumer, then take target throughput divided by the smaller of the two. For example, 60 MB/s peak with a consumer that handles 5 MB/s per partition needs at least 12. Also make sure the count is at least the largest number of consumers you expect in a group, and leave headroom for growth.
More partitions are not free: more files and open handles, more replication work, longer leader elections and more metadata. KRaft raises the cluster limits substantially, but tens of thousands of tiny partitions still cost.
You can increase the count with kafka-topics.sh --alter --partitions, but never decrease it. Increasing remaps keys (hash mod N changes), so per-key ordering breaks across the change, and existing data stays where it is. For keyed topics, choose generously up front or migrate to a new topic.
Likely follow-up: How would you migrate a keyed topic from 6 to 24 partitions without breaking per-key order?
CommitFailedException. How do you debug it?hardThe pattern points to members exceeding max.poll.interval.ms (5 minutes by default). Heartbeats run on a background thread, but if the application does not call poll() again within that interval, the consumer leaves the group. Its partitions move, its later commit fails with CommitFailedException, the records are reprocessed elsewhere, and the cycle repeats.
Confirm it from logs (a message that the poll timeout has expired), rebalance metrics and per-record processing time. A common trigger is one slow dependency multiplied by max.poll.records (500 by default).
Fixes, in order:
max.poll.records so a batch finishes well within the interval.max.poll.interval.ms if long batches are legitimate.pause() the partitions, process elsewhere, keep polling, then resume().Also check for session timeouts from long GC pauses, which look similar.
Likely follow-up: How is max.poll.interval.ms different from session.timeout.ms?
Go hop by hop and check each setting:
acks=0 or acks=1; send failures ignored because nobody checks the Future or callback; delivery.timeout.ms expiring during an outage; records still in the buffer when the process exits without flush() or close().min.insync.replicas=1 so acks=all meant only the leader; unclean.leader.election.enable=true; retention shorter than the time a consumer was down.auto.offset.reset=latest on a new group or after an out-of-range offset; exceptions caught and swallowed; a DLQ nobody watches.Then prove it: compare produced counts per partition (kafka-get-offsets.sh) with consumed counts, trace a missing event ID through producer logs, and check broker logs for unclean elections or ISR shrinks. The durable baseline is RF 3, min.insync.replicas=2, acks=all with idempotence, checked callbacks, and commit after processing.
Likely follow-up: How would you detect silent loss continuously instead of after a complaint?
No questions match that filter.
Prefer multiple choice? All 20 Apache Kafka MCQs with answers →