You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
{{ message }}
Repository navigation
Revert "fix(kafka): stop retrying lost consumer membership" - #1021
I saw this and pointed Claude at both franz-go and file.d. There's little to nothing that changed in the autocommit path between the versions that are changing in this PR. MarkCommitOffsets, the revoke path, the autocommit loop, and how record epochs are set are all unchanged. The commit-path changes are: a new commit waits for the previous one to finish instead of canceling it, which keeps commits in order; a classic group rejoins right away if a commit returns UNKNOWN_MEMBER_ID or ILLEGAL_GENERATION; and OffsetCommit v10 support, which drops back to v9 if any topic has no topic ID.
In the file.d audit, in plugin/input/kafka, it did find a few problems that can look like autocommit not working (my own rewriting of its response below):
client.go sets OnPartitionsRevoked(s.Lost). That replaces the default revoke, and the default revoke is where kgo does a blocking commit. Lost does not commit. Currently, on every rebalance, anything marked since the last autocommit is never committed, up to auto_commit_interval, and the next partition owner reprocesses it. Calling cl.CommitMarkedOffsets(ctx) at the start of the revoke callback fixes this.
Plugin.Commit calls MarkCommitOffsets whenever the pipeline finishes an event, even for partitions this member no longer owns. kgo clears its uncommitted state on revoke, but marking again adds the partition back, and then the next autocommit commits it. Classic groups accept that commit, so it can overwrite the new owner's committed offset (this may be irrelevant if the new partition owner does another commit over the stale-partition-owner's commit, but there is risk).
assembleOffset packs the offset and leader epoch into one int64, with the epoch in the low 16 bits. A record's leader epoch is -1 if its batch has no epoch (older message formats). offset<<16 + (-1) decodes as offset-1 with epoch 65535, so file.d commits the record's own offset instead of offset+1, with epoch 65535. The next time the partition is assigned, kgo sees a committed epoch and validates it with OffsetForLeaderEpoch. No partition has epoch 65535, so the broker answers with an undefined end offset. kgo reports that as data loss and resets the partition to ConsumeResetOffset. file.d's default offset is newest, so the consumer skips everything produced since the last commit; with oldest, it re-reads the whole partition. (An epoch of 65536 or higher also breaks the packing, but it's unlikely an epoch will get that high.)
Lost waits for each partition goroutine to exit. A goroutine blocked in controller.In cannot exit until the pipeline's event pool has room. consume() can block sending to a full fetches channel before it calls AllowRebalance. With BlockRebalanceOnPoll set, either of these can delay a rebalance past the rebalance timeout (60s by default). The broker then removes the member, and the member rejoins with a new member ID after Lost returns, which causes another rebalance. This is the same on 1.20.7 and 1.21.3: kgo keeps heartbeating while the revoke callback runs, so both versions see UNKNOWN_MEMBER_ID from the heartbeat, and the rejoin that 1.21.3 queues from a failed commit is dropped when the join starts. Each of these rebalances also triggers problems 1 and 2 above.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Reverts #1017
suspicions that autocommit is broken in the new versions franz-go