New to Kafka? Start with How Kafka Moves a Message: topics, partitions, offsets, replication and the basic consumer group. This page zooms in on the group.
The idea: a group shares partitions, and offsets decide what happens after a failure
A consumer group is a set of consumers with the same group.id. Each partition is read by exactly one member, so the members split the work, and the group as a whole reads every record. On the page, group billing reads two topics, orders (o0, o1, o2) and payments (p0, p1, p2), and turns each record into an invoice in a third topic, invoices. A record is named by partition and offset: o2.3 is offset 3 of orders-2, and its invoice is inv(o2.3).
Two questions make up the whole page. Who reads which partition, and what does it cost to change that? That is rebalancing. And where does the next owner of a partition start? That is the committed offset, and the gap between processing a record and committing it is where duplicates and lost records come from.
The group coordinator and __consumer_offsets
Every group has a group coordinator: the broker that leads the partition of the internal topic __consumer_offsets that the group id hashes to (B1 on the page). A consumer finds it with FindCoordinator. The coordinator keeps the member list and the generation (a number that grows with every rebalance), and it stores the committed offsets: an OffsetCommit is just a record written to __consumer_offsets. A committed offset is the next record to read.
The classic protocol: JoinGroup, SyncGroup, generations
When the membership changes, the coordinator moves to PreparingRebalance and answers the members' next heartbeat with REBALANCE_IN_PROGRESS. Then:
- JoinGroup. Every member sends
JoinGroupwith its subscription (and, with the cooperative protocol, the partitions it owns). The coordinator waits until every known member has joined, at most the rebalance timeout (the members'max.poll.interval.ms); a member that does not come is removed. - A new generation and a leader. The coordinator increments the generation and answers every
JoinGroup. One member is the group leader (the current leader if it is still there, otherwise the first to join), and only the leader gets the full member list. - SyncGroup. The leader runs the assignor on its own machine and sends the plan in its
SyncGroup; the others send an emptySyncGroup. The coordinator answers each member with its share and the group isStable. The broker never looks inside the plan: the assignment strategy is client code. - OffsetFetch. Each member reads the committed offsets of its new partitions and starts fetching from there.
Every offset commit carries the generation and the member id. A commit from an old generation, or from a member that is no longer in the group, is rejected with ILLEGAL_GENERATION or UNKNOWN_MEMBER_ID (the client throws CommitFailedException). This is how the group protects itself from a member that has lost its partitions but does not know it yet.
Two timeouts: session.timeout.ms and max.poll.interval.ms
A consumer has a background heartbeat thread. If the coordinator hears nothing for session.timeout.ms (45 s by default; 6 ticks on the page), the member is declared dead and the group rebalances. That catches crashes and lost machines, but a process that is alive and stuck (a long GC pause, a slow database call, an endless loop) keeps heartbeating. So there is a second limit: if the application does not call poll() again within max.poll.interval.ms (5 min by default; 12 ticks on the page), the heartbeat thread itself leaves the group. The process is still running, though, and when it finally finishes its batch it tries to commit: it is a zombie, and its commit is rejected (Demo 7). The two timeouts exist so that crash detection can be fast while slow processing is allowed.
A graceful close() sends LeaveGroup, so the group rebalances at once instead of waiting for the session timeout (Demo 6).
Assignors: range, round-robin, sticky, cooperative-sticky
| Assignor | How it splits | Protocol |
|---|---|---|
| range (the old default) | per topic: sort the partitions and the members, give each member a contiguous range, the first members one more | eager |
| round-robin | all partitions of all topics, dealt out one by one | eager |
| sticky | balanced, and moves as few partitions as possible | eager |
| cooperative-sticky | like sticky | cooperative (incremental) |
Because range works per topic, two topics of 3 partitions and 2 members give 4 / 2, not 3 / 3: each topic's extra partition goes to the first member. With 4 members, each topic has only 3 partitions to hand out, so C4 gets nothing although there are 6 partitions (Demo 2). Round-robin would fix the balance; the page keeps the choice to range vs cooperative-sticky and leaves round-robin to this paragraph.
| Members | range (eager) | cooperative-sticky |
|---|---|---|
| C1 | C1 all 6 | C1 all 6 |
| + C2 | C1 {o0 o1 p0 p1} · C2 {o2 p2}; all 6 paused | C1 {o0 o1 o2} · C2 {p0 p1 p2}; 3 paused |
| + C3 | C1 {o0 p0} · C2 {o1 p1} · C3 {o2 p2}; all 6 paused | C1 {o0 o1} · C2 {p0 p1} · C3 {o2 p2}; 2 paused |
| + C4 | the same three, C4 idle; all 6 paused | C1 {o0} · C2 {p0 p1} · C3 {o2 p2} · C4 {o1}; 1 paused |
Eager vs cooperative rebalancing
With the eager protocol every member revokes all its partitions before it sends JoinGroup (committing its position in onPartitionsRevoked), so for the whole rebalance nobody consumes anything: "stop the world". A member that is in the middle of a batch finishes it first, and the whole group waits for it.
With the cooperative (incremental) protocol (KIP-429), members keep their partitions and keep consuming during the rebalance; JoinGroup tells the leader what each member owns. The leader computes the target, but a partition that another member still owns is not handed out yet. After SyncGroup each member revokes only the partitions that are missing from its new assignment and rejoins at once; the second rebalance hands those partitions to their new owners. Two short rebalances, and only the partitions that move ever stop (Demo 3; compare the pause meter with Demo 2). The page's pause meter adds up paused partitions × ticks per rebalance.
Static membership
A rolling restart of a dynamic group costs two rebalances per instance: one when it leaves, one when it comes back. With group.instance.id (KIP-345) a member is static: it does not send LeaveGroup on close, and when an instance with the same id comes back within session.timeout.ms, the coordinator swaps in the new member id and returns the same assignment, with no rebalance and no new generation (Demo 6). The price is that a static member that really died is only noticed after the session timeout, which is why static groups usually raise it.
The new consumer group protocol (KIP-848)
Kafka 4.0 made a new protocol generally available (group.protocol=consumer). The coordinator computes the assignment, members send one kind of request (ConsumerGroupHeartbeat) and each member moves to its target at its own pace: there is no global barrier where everyone waits for the slowest member, and no client-side leader. Revocations are still incremental, and offsets, generations (member epochs) and fencing work the same way. The page draws the classic protocol, which most running clusters still use and whose phases are easier to see.
Offsets and delivery semantics
The next owner of a partition starts at the committed offset, not where the previous owner stopped. So the order of processing and committing decides what a failure does:
- Commit after processing (
commitSync()after the batch): a crash between the two means the batch is processed again. This is at-least-once, and it is where duplicates come from (Demo 4). - Commit before processing: a crash in between skips the rest of the batch. This is at-most-once (Demo 5:
p2.3is never processed). - Auto-commit (
enable.auto.commit, the default):poll()commits the positions of earlier polls everyauto.commit.interval.ms. Since it only commits what earlier polls returned, a single-threaded loop is at-least-once, with more duplicates after a crash than commit-after. - A zombie (a member thrown out while it was still processing) finishes its batch, its commit is rejected, and the new owner processes the same records again (Demo 7).
With at-least-once the usual fix is an idempotent sink: an upsert by record id, a unique key, a "processed ids" table. When the output goes to another Kafka topic, Kafka can do better.
The idempotent producer: PID, sequence numbers, epoch
With enable.idempotence=true (the default since Kafka 3.0) the producer gets a producer id (PID) from InitProducerId and numbers its records per partition with a sequence number. The partition leader remembers the last 5 sequences of each PID; a record that arrives again with a sequence it already has is not written a second time, and the answer carries the original offset. That removes the duplicates of retries: a produce whose answer was lost (Demo 8).
It does not remove duplicates of re-processing. When C2 crashes and C1 processes o2.3 again, C1 sends inv(o2.3) with its own PID and a new sequence number: for the broker it is a new record (Demo 9). Sequence numbers only protect one producer session.
Transactions: output and offsets commit together
A consume-process-produce loop is exactly-once if the output records and the consumer's offsets are committed atomically. That is what Kafka transactions (KIP-98) do. The producer has a transactional.id (tx-C1, …). For each batch:
InitProducerId(transactional.id)once at start: the transaction coordinator (the leader of the__transaction_statepartition for that id; B1 on the page) returns the PID and a bumped epoch, and aborts any transaction a previous instance left open.beginTransaction(); the first send to a partition adds it:AddPartitionsToTxn(invoices-0). The coordinator logsOngoing.- The output records are written to
invoices-0as usual, marked as transactional (white cells on the page). sendOffsetsToTransaction():AddOffsetsToTxn(billing)to the transaction coordinator, thenTxnOffsetCommitto the group coordinator. The offsets are stored as pending: nobody sees them yet.commitTransaction():EndTxn(COMMIT). The coordinator logsPrepareCommit(from here on the outcome is decided), sendsWriteTxnMarkersto every partition in the transaction, and logsCompleteCommit. A COMMIT marker ininvoices-0makes the invoices visible, and the marker in__consumer_offsetsmakes the offsets committed, in one step.
A transaction that stays open longer than transaction.timeout.ms (60 s by default; 30 ticks on the page) is aborted by the coordinator: an ABORT marker, and the records are dropped by readers.
Readers: read_committed, the LSO and aborted records
Transactional records are written to the log before anyone knows whether they will commit. A consumer with isolation.level=read_committed reads only up to the last stable offset (LSO): the first offset of the oldest transaction that is still open. It skips records of aborted transactions using the partition's aborted transaction index. So one open transaction holds back every reader, even for records of other producers that committed long ago (Demo 11: after C2's crash, the audit reader sees nothing past inv(o2.3) until the transaction times out). read_uncommitted (the default!) reads up to the high watermark and sees aborted records too; switch the audit reader to it at the end of Demo 11 to see the duplicates come back.
Zombie fencing: producer epoch and group generation
A zombie must not be able to commit anything after its replacement took over. Two fences do that:
- The producer epoch. The new instance's
InitProducerIdwith the sametransactional.idbumps the epoch. Every request of the old instance carries the old epoch and is refused withPRODUCER_FENCED; its open transaction is aborted (Demo 12, second part; Demo 11's variant with Restart right after the crash, which aborts at once instead of waiting for the timeout). - The group generation (KIP-447, Kafka 2.5).
TxnOffsetCommitcarries the consumer's generation and member id, so a member that was thrown out of the group is refused withILLEGAL_GENERATIONeven when it has its owntransactional.id; it has to abort (Demo 12, first part). Before KIP-447 each input partition needed its own transactional producer so that the epoch alone could fence zombies.
A new owner also must not start from offsets that a transaction is still committing: OffsetFetch answers UNSTABLE_OFFSET_COMMIT while a pending transactional offset exists for the partition, and the consumer asks again.
Exactly-once has limits
Exactly-once covers reading from Kafka, writing to Kafka and committing offsets to Kafka. An external side effect (a database write, an HTTP call, an e-mail) is not in the transaction: if the consumer crashes after it, it will happen again. Make the side effect idempotent, or write it to a Kafka topic first and apply it from there (the outbox pattern). A distributed transaction across Kafka and a database is the two-phase commit problem, see Saga vs 2PC.
Kafka Streams does all of this for you with processing.guarantee=exactly_once_v2: transactional producers, offsets in the transaction, read_committed input and fencing by generation.
How the page simplifies
- Replication factor 1 and brokers never fail (Kafka.html covers replication,
acksand leader failover); B1 is both the group coordinator and the transaction coordinator. - Heartbeats are not drawn: the coordinator hears from every live member every tick and removes one after exactly 6 ticks of silence. The rebalance notification is drawn as its own packet instead of a heartbeat answer.
- A poll is one fetch round trip for all of a member's partitions,
max.poll.records= 3, processing takes 1 tick per record (15 for a "slow" record); every message takes 1 tick. The timeouts are small tick counts: session 6, max.poll.interval 12, auto-commit 5, transaction 30, producer retry 3. - The cooperative-sticky assignor is a small deterministic version of the real sticky rules: keep what each member owns, give unowned partitions in order to the member with the fewest (ties: lower id), then while the counts differ by more than one, the member with the most gives its last partition to the member with the fewest. The real one can pick different partitions.
- Changing the assignor takes effect at the next rebalance, without the rolling upgrade a real switch from eager to cooperative needs.
- One output partition, and each poll batch is one transaction; real applications often commit every N records or milliseconds. With a plain or idempotent producer the page flushes the producer before a commit-after.
- An open transaction's records are white; aborted ones are red ✗ cells that
read_committedskips (the aborted transaction index itself is not drawn). - Restart of a member stuck in a slow record leaves the old process running as a zombie on purpose (a stand-in for "the orchestrator started a replacement during a GC pause"). Crash (after 2 records) and the start of a slow record stop the clock, so that you can act before the group reacts.
- Left out: incremental fetch sessions, rack-aware assignment, the transaction coordinator's own failover, consumer lag metrics, pattern subscriptions,
isolation.levelof the group's own input, andMEMBER_ID_REQUIREDon a first JoinGroup.
References
Kafka documentation: message delivery semantics
KIP-98: Exactly Once Delivery and Transactional Messaging
KIP-429: Kafka Consumer Incremental Rebalance Protocol
KIP-447: Producer scalability for exactly once semantics (fencing by generation)
KIP-848: The Next Generation of the Consumer Rebalance Protocol