KAFKA-13185: Clear pending records after pre-commit rewind - #23048
Open
lh0156 wants to merge 1 commit into
Open
Conversation
Clear the pending message batch and its associated offsets when a sink task pre-commit fails and the consumer is rewound to the last committed position. Add a regression test covering a retriable put followed by a failed pre-commit. Generated-by: OpenAI Codex (GPT-5) Signed-off-by: Yunseop Eom <62834176+lh0156@users.noreply.github.com>
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
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
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.
Closes KAFKA-13185
Summary
messageBatchwhen a sink task pre-commit fails and the consumer is rewound to the last committed position.origOffsetsas well, so the discarded batch cannot advancecurrentOffsetsduring the recovery poll.put, a failedpreCommit, and the subsequent empty delivery.When
SinkTask.putraises aRetriableException,WorkerSinkTaskkeeps the batch for redelivery. If the followingpreCommitalso fails, the worker seeks to the last committed offset but previously retained that batch. The next poll could therefore deliver stale records, and retainingorigOffsetscould reintroduce offsets that had just been rewound. The recovery path now discards both pieces of pending state before polling again.Tests
./gradlew :connect:runtime:test --tests org.apache.kafka.connect.runtime.WorkerSinkTaskTest.testPreCommitFailureClearsPendingMessageBatch --no-build-cache --console=plain./gradlew :connect:runtime:spotlessCheck :connect:runtime:test --tests org.apache.kafka.connect.runtime.WorkerSinkTaskTest --no-build-cache --console=plain./gradlew :connect:runtime:test --no-build-cache --console=plainThe regression test was verified RED before the fix because the stale batch was delivered with one record, then GREEN after the fix with an empty delivery and last-committed offsets preserved.