Repository navigation
[fix][client] Send the pending batch before a non-batched message so deduplication does not drop it - #616
Open
gusteycamargo wants to merge 1 commit into
Conversation
…deduplication does not drop it Fixes apache#615 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.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.
Fixes #615
Motivation
When batching is enabled, a message that cannot be batched (e.g. a delayed message) is sent immediately, while the messages already in the batch container still wait for the batch to be flushed. The non-batched message took the next sequence id, so it reaches the broker before messages with lower sequence ids. With deduplication enabled, the broker then drops the whole batch as a duplicate, while the producer receives
ResultOkfor every message.The Java client avoids it by flushing the batch container before sending any non-batched message (
ProducerImpl#processOpSendMsg).Modifications
ProducerImpl::sendAsync, before sending a message that cannot be added to the batch, send the pending batch (batchMessageAndSend().complete()), the same way the producer already does when the batch container has no room for the next message.Verifying this change
This change added tests and can be verified as follows:
ClientDeduplicationBatchingTest.testBatchedMessagesBeforeDelayedMessage, parameterized forDefaultBatchingandKeyBasedBatching: with deduplication enabled on the namespace, it sends 10 batched messages asynchronously and then a delayed one, and reads the topic from the earliest position. Without the fix only 1 of the 11 messages is stored, with it all 11 are.Documentation
doc-required(Your PR needs to update docs and you will update later)
doc-not-needed(Bug fix in the producer, no API change)
🤖 Generated with Claude Code