-
Notifications
You must be signed in to change notification settings - Fork 74
[FLINK-35477] Make initial cursor position configurable #105
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Conversation
|
Thanks for opening this pull request! Please check out our contributing guidelines. (https://flink.apache.org/contributing/how-to-contribute.html) |
|
In Pull Request #103, only the |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Logic LGTM, added 2 minor suggestions.
| } | ||
| } | ||
|
|
||
| if (Objects.equals(startCursor, StartCursor.latest())) { |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
nit: StartCursor.latest().equals(startCursor) is also null-safe, shorter and does not require a util class
| .enumType(SubscriptionInitialPosition.class) | ||
| .defaultValue(SubscriptionInitialPosition.Latest) | ||
| .withDescription( | ||
| Description.builder().text("Consumer initial position.").build()); |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I would add a bit more descriptive text here, e.g.: Initial cursor position of the consumer.
Purpose of the change
For example: Add dynamic sink topic support for Pulsar connector.
Brief change log
ProducerRegister.PulsarSinkContext.MetadataListener.Verifying this change
Please make sure both new and modified tests in this PR follows the conventions defined in our code quality
guide: https://flink.apache.org/contributing/code-style-and-quality-common.html#testing
(Please pick either of the following options)
This change is a trivial rework / code cleanup without any test coverage.
(or)
This change is already covered by existing tests, such as (please describe tests).
(or)
This change added tests and can be verified as follows:
(example:)
Significant changes
(Please check any boxes [x] if the answer is "yes". You can first publish the PR and check them afterwards, for
convenience.)
@Public(Evolving))