Skip to content

Commit bc7f3fb

Browse files
authored
Bugfix: Python SDK listing of ingestion job fails for featureset reference filter (#861)
* Fix listing of ingestion jobs for featureset ref * Fix lint * Shift filter to java method
1 parent beed7e5 commit bc7f3fb

3 files changed

Lines changed: 20 additions & 17 deletions

File tree

core/src/main/java/feast/core/dao/JobRepository.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,9 @@ public interface JobRepository extends JpaRepository<Job, String> {
4040

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

43+
List<Job> findByFeatureSetJobStatusesFeatureSetNameAndFeatureSetJobStatusesFeatureSetProjectName(
44+
String featureSetName, String featureSetProject);
45+
4346
// find jobs that have at least one store with given name
4447
List<Job> findByJobStoresIdStoreName(String storeName);
4548
}

core/src/main/java/feast/core/service/JobService.java

Lines changed: 8 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -23,11 +23,9 @@
2323
import feast.core.log.Action;
2424
import feast.core.log.AuditLogger;
2525
import feast.core.log.Resource;
26-
import feast.core.model.FeatureSet;
2726
import feast.core.model.Job;
2827
import feast.core.model.JobStatus;
2928
import feast.proto.core.CoreServiceProto.ListFeatureSetsRequest;
30-
import feast.proto.core.CoreServiceProto.ListFeatureSetsResponse;
3129
import feast.proto.core.CoreServiceProto.ListIngestionJobsRequest;
3230
import feast.proto.core.CoreServiceProto.ListIngestionJobsResponse;
3331
import feast.proto.core.CoreServiceProto.RestartIngestionJobRequest;
@@ -115,20 +113,17 @@ public ListIngestionJobsResponse listJobs(ListIngestionJobsRequest request)
115113
if (filter.hasFeatureSetReference()) {
116114
// find a matching featuresets for reference
117115
FeatureSetReference fsReference = filter.getFeatureSetReference();
118-
ListFeatureSetsResponse response =
119-
this.specService.listFeatureSets(this.toListFeatureSetFilter(fsReference));
120-
List<FeatureSet> featureSets =
121-
response.getFeatureSetsList().stream()
122-
.map(FeatureSet::fromProto)
123-
.collect(Collectors.toList());
124116

125117
// find jobs for the matching featuresets
126118
Collection<Job> matchingJobs =
127-
this.jobRepository.findByFeatureSetJobStatusesIn(
128-
featureSets.stream()
129-
.flatMap(fs -> fs.getJobStatuses().stream())
130-
.collect(Collectors.toList()));
131-
List<String> jobIds = matchingJobs.stream().map(Job::getId).collect(Collectors.toList());
119+
this.jobRepository
120+
.findByFeatureSetJobStatusesFeatureSetNameAndFeatureSetJobStatusesFeatureSetProjectName(
121+
fsReference.getName(), fsReference.getProject());
122+
List<String> jobIds =
123+
matchingJobs.stream()
124+
.filter(job -> job.getStatus().equals(JobStatus.RUNNING))
125+
.map(Job::getId)
126+
.collect(Collectors.toList());
132127
matchingJobIds = this.mergeResults(matchingJobIds, jobIds);
133128
}
134129
}

sdk/python/feast/client.py

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -749,16 +749,21 @@ def list_ingest_jobs(
749749
List of IngestJobs matching the given filters
750750
"""
751751
# construct list request
752-
feature_set_ref = None
752+
feature_set_ref_proto = None
753+
if feature_set_ref:
754+
feature_set_ref_proto = feature_set_ref.to_proto()
753755
list_filter = ListIngestionJobsRequest.Filter(
754-
id=job_id, feature_set_reference=feature_set_ref, store_name=store_name,
756+
id=job_id,
757+
feature_set_reference=feature_set_ref_proto,
758+
store_name=store_name,
755759
)
756760
request = ListIngestionJobsRequest(filter=list_filter)
757761
# make list request & unpack response
758-
response = self._core_service_stub.ListIngestionJobs(request) # type: ignore
762+
response = self._core_service.ListIngestionJobs(request) # type: ignore
759763
ingest_jobs = [
760-
IngestJob(proto, self._core_service_stub) for proto in response.jobs # type: ignore
764+
IngestJob(proto, self._core_service) for proto in response.jobs # type: ignore
761765
]
766+
762767
return ingest_jobs
763768

764769
def restart_ingest_job(self, job: IngestJob):

0 commit comments

Comments
 (0)