Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions core/src/main/java/feast/core/dao/JobRepository.java
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,9 @@ public interface JobRepository extends JpaRepository<Job, String> {

List<Job> findByFeatureSetJobStatusesIn(List<FeatureSetJobStatus> featureSetsJobStatuses);

List<Job> findByFeatureSetJobStatusesFeatureSetNameAndFeatureSetJobStatusesFeatureSetProjectName(
String featureSetName, String featureSetProject);

// find jobs that have at least one store with given name
List<Job> findByJobStoresIdStoreName(String storeName);
}
21 changes: 8 additions & 13 deletions core/src/main/java/feast/core/service/JobService.java
Original file line number Diff line number Diff line change
Expand Up @@ -23,11 +23,9 @@
import feast.core.log.Action;
import feast.core.log.AuditLogger;
import feast.core.log.Resource;
import feast.core.model.FeatureSet;
import feast.core.model.Job;
import feast.core.model.JobStatus;
import feast.proto.core.CoreServiceProto.ListFeatureSetsRequest;
import feast.proto.core.CoreServiceProto.ListFeatureSetsResponse;
import feast.proto.core.CoreServiceProto.ListIngestionJobsRequest;
import feast.proto.core.CoreServiceProto.ListIngestionJobsResponse;
import feast.proto.core.CoreServiceProto.RestartIngestionJobRequest;
Expand Down Expand Up @@ -115,20 +113,17 @@ public ListIngestionJobsResponse listJobs(ListIngestionJobsRequest request)
if (filter.hasFeatureSetReference()) {
// find a matching featuresets for reference
FeatureSetReference fsReference = filter.getFeatureSetReference();
ListFeatureSetsResponse response =
this.specService.listFeatureSets(this.toListFeatureSetFilter(fsReference));
List<FeatureSet> featureSets =
response.getFeatureSetsList().stream()
.map(FeatureSet::fromProto)
.collect(Collectors.toList());

// find jobs for the matching featuresets
Collection<Job> matchingJobs =
this.jobRepository.findByFeatureSetJobStatusesIn(
featureSets.stream()
.flatMap(fs -> fs.getJobStatuses().stream())
.collect(Collectors.toList()));
List<String> jobIds = matchingJobs.stream().map(Job::getId).collect(Collectors.toList());
this.jobRepository
.findByFeatureSetJobStatusesFeatureSetNameAndFeatureSetJobStatusesFeatureSetProjectName(
fsReference.getName(), fsReference.getProject());
List<String> jobIds =
matchingJobs.stream()
.filter(job -> job.getStatus().equals(JobStatus.RUNNING))
.map(Job::getId)
.collect(Collectors.toList());
matchingJobIds = this.mergeResults(matchingJobIds, jobIds);
}
}
Expand Down
13 changes: 9 additions & 4 deletions sdk/python/feast/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -749,16 +749,21 @@ def list_ingest_jobs(
List of IngestJobs matching the given filters
"""
# construct list request
feature_set_ref = None
feature_set_ref_proto = None
if feature_set_ref:
feature_set_ref_proto = feature_set_ref.to_proto()
list_filter = ListIngestionJobsRequest.Filter(
id=job_id, feature_set_reference=feature_set_ref, store_name=store_name,
id=job_id,
feature_set_reference=feature_set_ref_proto,
store_name=store_name,
)
request = ListIngestionJobsRequest(filter=list_filter)
# make list request & unpack response
response = self._core_service_stub.ListIngestionJobs(request) # type: ignore
response = self._core_service.ListIngestionJobs(request) # type: ignore
ingest_jobs = [
IngestJob(proto, self._core_service_stub) for proto in response.jobs # type: ignore
IngestJob(proto, self._core_service) for proto in response.jobs # type: ignore
]

return ingest_jobs

def restart_ingest_job(self, job: IngestJob):
Expand Down