-
Notifications
You must be signed in to change notification settings - Fork 1k
New issue
Have a question about this project? # for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “#”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? # to your account
Add configurable prefix to Consumer Group in IngestionJob's Kafka reader #969
Conversation
/retest |
fe37fa6
to
cde8bbe
Compare
/retest |
1 similar comment
/retest |
@@ -23,6 +23,9 @@ feast: | |||
# Enabling JobManagement | |||
enabled: true | |||
|
|||
# Prefix for JobId |
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.
This comment says absolutely nothing.
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.
Updated comment.
@@ -71,6 +74,9 @@ private String createJobId(SourceProto.Source source) { | |||
source.getKafkaSourceConfig().getBootstrapServers(), | |||
source.getKafkaSourceConfig().getTopic()), | |||
dateSuffix); | |||
if (!this.jobProperties.getJobIdPrefix().isEmpty()) { |
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.
What happens when jobProperties is null?
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.
Addressed case.
@@ -87,9 +87,9 @@ public JobGroupingStrategy getJobGroupingStrategy( | |||
Boolean shouldConsolidateJobs = | |||
feastProperties.getJobs().getController().getConsolidateJobsPerSource(); | |||
if (shouldConsolidateJobs) { | |||
return new ConsolidatedJobStrategy(jobRepository); | |||
return new ConsolidatedJobStrategy(jobRepository, feastProperties.getJobs()); |
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.
Can we extract the jobProperties before passing it? It's a convention that I am trying to get us to follow, even if it takes more lines of code
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.
Extracted before passing jobProperties.
[APPROVALNOTIFIER] This PR is APPROVED This pull-request has been approved by: pyalex, terryyylim The full list of commands accepted by this bot can be found here. The pull request process is described here
Needs approval from an approver in each of these files:
Approvers can indicate their approval by writing |
/lgtm |
What this PR does / why we need it:
To allow multiple deployments with jobs running in parallel, we should separate their Kafka consumer groups which is based on the JobId. This would prevent overwriting each other's offsets.
Which issue(s) this PR fixes:
Fixes #
Does this PR introduce a user-facing change?: