Skip to content

Commit 4105a31

Browse files
committed
Add JobId prefix setting
1 parent 07625dd commit 4105a31

5 files changed

Lines changed: 23 additions & 4 deletions

File tree

job-controller/src/main/java/feast/jobcontroller/config/FeastProperties.java

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -176,6 +176,10 @@ public feast.jobcontroller.runner.Runner getType() {
176176
/* Population job metric properties */
177177
private MetricsProperties metrics;
178178

179+
@NotNull
180+
/* Prefix for JobId to separate consumer groups for independent jobs running in parallel */
181+
private String jobIdPrefix;
182+
179183
/* Timeout in seconds for each attempt to update or submit a new job to the runner */
180184
@Positive private long jobUpdateTimeoutSeconds;
181185

job-controller/src/main/java/feast/jobcontroller/config/JobControllerConfig.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -87,9 +87,9 @@ public JobGroupingStrategy getJobGroupingStrategy(
8787
Boolean shouldConsolidateJobs =
8888
feastProperties.getJobs().getController().getConsolidateJobsPerSource();
8989
if (shouldConsolidateJobs) {
90-
return new ConsolidatedJobStrategy(jobRepository);
90+
return new ConsolidatedJobStrategy(jobRepository, feastProperties.getJobs());
9191
} else {
92-
return new JobPerStoreStrategy(jobRepository);
92+
return new JobPerStoreStrategy(jobRepository, feastProperties.getJobs());
9393
}
9494
}
9595

job-controller/src/main/java/feast/jobcontroller/runner/ConsolidatedJobStrategy.java

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
*/
1717
package feast.jobcontroller.runner;
1818

19+
import feast.jobcontroller.config.FeastProperties.JobProperties;
1920
import feast.jobcontroller.dao.JobRepository;
2021
import feast.jobcontroller.model.Job;
2122
import feast.jobcontroller.model.JobStatus;
@@ -38,9 +39,11 @@
3839
*/
3940
public class ConsolidatedJobStrategy implements JobGroupingStrategy {
4041
private final JobRepository jobRepository;
42+
private final JobProperties jobProperties;
4143

42-
public ConsolidatedJobStrategy(JobRepository jobRepository) {
44+
public ConsolidatedJobStrategy(JobRepository jobRepository, JobProperties jobProperties) {
4345
this.jobRepository = jobRepository;
46+
this.jobProperties = jobProperties;
4447
}
4548

4649
@Override
@@ -71,6 +74,9 @@ private String createJobId(SourceProto.Source source) {
7174
source.getKafkaSourceConfig().getBootstrapServers(),
7275
source.getKafkaSourceConfig().getTopic()),
7376
dateSuffix);
77+
if (!this.jobProperties.getJobIdPrefix().isEmpty()) {
78+
jobId = this.jobProperties.getJobIdPrefix() + "-" + jobId;
79+
}
7480
return jobId.replaceAll("_store", "-").toLowerCase();
7581
}
7682

job-controller/src/main/java/feast/jobcontroller/runner/JobPerStoreStrategy.java

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717
package feast.jobcontroller.runner;
1818

1919
import com.google.common.collect.Lists;
20+
import feast.jobcontroller.config.FeastProperties.JobProperties;
2021
import feast.jobcontroller.dao.JobRepository;
2122
import feast.jobcontroller.model.Job;
2223
import feast.jobcontroller.model.JobStatus;
@@ -38,9 +39,11 @@
3839
*/
3940
public class JobPerStoreStrategy implements JobGroupingStrategy {
4041
private final JobRepository jobRepository;
42+
private final JobProperties jobProperties;
4143

42-
public JobPerStoreStrategy(JobRepository jobRepository) {
44+
public JobPerStoreStrategy(JobRepository jobRepository, JobProperties jobProperties) {
4345
this.jobRepository = jobRepository;
46+
this.jobProperties = jobProperties;
4447
}
4548

4649
@Override
@@ -78,6 +81,9 @@ private String createJobId(SourceProto.Source source, Iterable<StoreProto.Store>
7881
source.getKafkaSourceConfig().getTopic()),
7982
Lists.newArrayList(stores).get(0).getName(),
8083
dateSuffix);
84+
if (!this.jobProperties.getJobIdPrefix().isEmpty()) {
85+
jobId = this.jobProperties.getJobIdPrefix() + "-" + jobId;
86+
}
8187
return jobId.replaceAll("_store", "-").toLowerCase();
8288
}
8389

job-controller/src/main/resources/application.yml

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,9 @@ feast:
2323
# Enabling JobManagement
2424
enabled: true
2525

26+
# Prefix for JobId
27+
job_id_prefix: ""
28+
2629
# Job update polling interval in milliseconds: how often Feast checks if new jobs should be sent to the runner.
2730
polling_interval_milliseconds: 60000
2831

0 commit comments

Comments
 (0)