fix: key pending commit work by consumer and partition - #99
Merged
Merged
Conversation
Pending lists were keyed by TopicPartition alone, so subscribers in different consumer groups reading the same topic shared one list per partition, and an unfinished task in one group held back the other group's ready prefix. Key them by (id(consumer), partition), the identity cancellation watermarks already use. Closes #96
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 #96.
Problem
PendingCommitskept one offset-sorted list perTopicPartition. When one concurrent handler serves subscribers in different consumer groups that read the same topic, both groups' user tasks for a partition shared that list, so the ready prefix was computed across both. An in-flight task in group A held back group B's finished offsets on that partition until A's task finished. Commits stayed safe, but one group's slow handler stalled the other group's commits and their share ofmax_uncommitted_tasks.Change
(id(consumer), partition), the same owner identity the cancellation watermarks already use. The newOwnerKeyalias names it once for both.extract_ready_prefixesis unchanged apart from its key type. The public surface ofPendingCommits(absorb,hold,take_ready,__len__,clear_watermarks) and the shape ofReadyCommit(still one per consumer) are unchanged.Tests
Regression tests written first. These fail on
main:test_another_consumers_unfinished_task_does_not_hold_back_a_ready_prefix(both orders)test_unfinished_task_still_holds_back_its_own_consumers_later_taskstest_another_consumer_groups_in_flight_task_does_not_delay_a_commit(committer-level, the reproduction from the issue)test_cancelled_task_sets_a_watermark_only_for_its_own_consumerpasses onmainand guards watermark scoping. The fourextract_ready_prefixestests now build pending state keyed byOwnerKey.just lintpasses.just test --cov=. --cov-branchgives 232 passed, 100% branch coverage.