[KafkaIO] Use consumer position and lag to estimate end offsets in KafkaUnboundedReader - #39830
sjvanrossum wants to merge 2 commits into
Conversation
|
Assigning reviewers: R: @kennknowles for label java. Note: If you would like to opt out of this review, comment Available commands:
The PR bot will only process comments in the main thread (not review comments). |
|
Reminder, please take a look at this pr: @kennknowles @johnjcasey |
|
Assigning new set of reviewers because Pr has gone too long without review. If you would like to opt out of this review, comment R: @Abacn for label java. Available commands:
|
|
Reminder, please take a look at this pr: @Abacn @kennknowles |
|
Stopping reviewer notifications for this pull request: review requested by someone other than the bot, ceding control. If you'd like to restart, comment |
Use
Consumer.position(TopicPartition)andConsumer.currentLag(TopicPartition)to update the latest offset inPartitionState.This change reuses local metadata from recently completed record fetches instead of periodically requesting the broker for end offsets.
This PR is a follow up to #39285 with additional changes to reduce synchronization overhead in
consumerPollLoopto a minimum sincePartitionStatewill likely be updated more frequently than the offset consumer interval.Thank you for your contribution! Follow this checklist to help us incorporate your contribution quickly and easily:
addresses #123), if applicable. This will automatically add a link to the pull request in the issue. If you would like the issue to automatically close on merging the pull request, commentfixes #<ISSUE NUMBER>instead.CHANGES.mdwith noteworthy changes.See the Contributor Guide for more tips on how to make review process smoother.
To check the build health, please visit https://github.com/apache/beam/blob/master/.test-infra/BUILD_STATUS.md
GitHub Actions Tests Status (on master branch)
See CI.md for more information about GitHub Actions CI or the workflows README to see a list of phrases to trigger workflows.