Skip to content

Commit 65626d6

Browse files
author
Oleksii Moskalenko
authored
JobCoordinator use public API to communicate with Core (#943)
* test job with no labels * filtering by status & status update in Core API * specs IT * SpecServiceIT * resolve conflicts * comments
1 parent 51f4fb6 commit 65626d6

14 files changed

Lines changed: 1035 additions & 1190 deletions

core/src/main/java/feast/core/grpc/CoreServiceImpl.java

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -167,6 +167,22 @@ public void getFeatureStatistics(
167167
}
168168
}
169169

170+
@Override
171+
public void updateFeatureSetStatus(
172+
UpdateFeatureSetStatusRequest request,
173+
StreamObserver<UpdateFeatureSetStatusResponse> responseObserver) {
174+
try {
175+
UpdateFeatureSetStatusResponse response = specService.updateFeatureSetStatus(request);
176+
177+
responseObserver.onNext(response);
178+
responseObserver.onCompleted();
179+
} catch (Exception e) {
180+
log.error("Exception has occurred in UpdateFeatureSetStatus method: ", e);
181+
responseObserver.onError(
182+
Status.INTERNAL.withDescription(e.getMessage()).withCause(e).asRuntimeException());
183+
}
184+
}
185+
170186
@Override
171187
public void listStores(
172188
ListStoresRequest request, StreamObserver<ListStoresResponse> responseObserver) {

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

Lines changed: 73 additions & 41 deletions
Original file line numberDiff line numberDiff line change
@@ -18,22 +18,24 @@
1818

1919
import static feast.common.models.Store.isSubscribedToFeatureSet;
2020
import static feast.core.model.FeatureSet.parseReference;
21-
import static feast.core.util.StreamUtil.wrapException;
2221

2322
import com.google.common.collect.Sets;
23+
import com.google.protobuf.InvalidProtocolBufferException;
2424
import feast.common.models.FeatureSetReference;
2525
import feast.core.config.FeastProperties;
2626
import feast.core.config.FeastProperties.JobProperties;
27-
import feast.core.dao.FeatureSetRepository;
2827
import feast.core.job.*;
2928
import feast.core.job.task.*;
30-
import feast.core.model.*;
31-
import feast.core.model.FeatureSet;
29+
import feast.core.model.FeatureSetDeliveryStatus;
3230
import feast.core.model.Job;
3331
import feast.core.model.JobStatus;
32+
import feast.proto.core.CoreServiceProto;
3433
import feast.proto.core.CoreServiceProto.ListStoresRequest.Filter;
3534
import feast.proto.core.CoreServiceProto.ListStoresResponse;
3635
import feast.proto.core.FeatureSetProto;
36+
import feast.proto.core.FeatureSetProto.FeatureSet;
37+
import feast.proto.core.FeatureSetProto.FeatureSetSpec;
38+
import feast.proto.core.FeatureSetReferenceProto;
3739
import feast.proto.core.IngestionJobProto;
3840
import feast.proto.core.SourceProto.Source;
3941
import feast.proto.core.StoreProto.Store;
@@ -59,24 +61,21 @@ public class JobCoordinatorService {
5961
private final int SPEC_PUBLISHING_TIMEOUT_SECONDS = 5;
6062

6163
private final JobRepository jobRepository;
62-
private final FeatureSetRepository featureSetRepository;
6364
private final SpecService specService;
6465
private final JobManager jobManager;
6566
private final JobProperties jobProperties;
6667
private final JobGroupingStrategy groupingStrategy;
67-
private final KafkaTemplate<String, FeatureSetProto.FeatureSetSpec> specPublisher;
68+
private final KafkaTemplate<String, FeatureSetSpec> specPublisher;
6869

6970
@Autowired
7071
public JobCoordinatorService(
7172
JobRepository jobRepository,
72-
FeatureSetRepository featureSetRepository,
7373
SpecService specService,
7474
JobManager jobManager,
7575
FeastProperties feastProperties,
7676
JobGroupingStrategy groupingStrategy,
77-
KafkaTemplate<String, FeatureSetProto.FeatureSetSpec> specPublisher) {
77+
KafkaTemplate<String, FeatureSetSpec> specPublisher) {
7878
this.jobRepository = jobRepository;
79-
this.featureSetRepository = featureSetRepository;
8079
this.specService = specService;
8180
this.jobManager = jobManager;
8281
this.jobProperties = feastProperties.getJobs();
@@ -236,7 +235,7 @@ private boolean jobRequiresUpgrade(Job job, Set<Store> stores) {
236235
*/
237236
FeatureSet allocateFeatureSetToJobs(FeatureSet featureSet) {
238237
FeatureSetReference ref =
239-
FeatureSetReference.of(featureSet.getProject().getName(), featureSet.getName());
238+
FeatureSetReference.of(featureSet.getSpec().getProject(), featureSet.getSpec().getName());
240239
Set<String> confirmedJobIds = new HashSet<>();
241240

242241
Stream<Pair<Source, Store>> jobArgsStream =
@@ -245,9 +244,9 @@ FeatureSet allocateFeatureSetToJobs(FeatureSet featureSet) {
245244
s ->
246245
isSubscribedToFeatureSet(
247246
s.getSubscriptionsList(),
248-
featureSet.getProject().getName(),
249-
featureSet.getName()))
250-
.map(s -> Pair.of(featureSet.getSource().toProto(), s));
247+
featureSet.getSpec().getProject(),
248+
featureSet.getSpec().getName()))
249+
.map(s -> Pair.of(featureSet.getSpec().getSource(), s));
251250

252251
// Add featureSet to allocated job if not allocated before
253252
for (Pair<Source, Set<Store>> jobArgs : groupingStrategy.collectSingleJobInput(jobArgsStream)) {
@@ -320,31 +319,47 @@ Iterable<Pair<Source, Set<Store>>> getSourceToStoreMappings() {
320319
* @param store to get subscribed FeatureSets for
321320
* @return list of FeatureSets that the store subscribes to.
322321
*/
323-
private List<FeatureSetProto.FeatureSet> getFeatureSetsForStore(Store store) {
322+
private List<FeatureSet> getFeatureSetsForStore(Store store) {
324323
return store.getSubscriptionsList().stream()
325324
.flatMap(
326-
subscription ->
327-
featureSetRepository
328-
.findAllByNameLikeAndProject_NameLikeOrderByNameAsc(
329-
subscription.getName().replace('*', '%'),
330-
subscription.getProject().replace('*', '%'))
331-
.stream())
325+
subscription -> {
326+
try {
327+
return specService
328+
.listFeatureSets(
329+
CoreServiceProto.ListFeatureSetsRequest.Filter.newBuilder()
330+
.setProject(subscription.getProject())
331+
.setFeatureSetName(subscription.getName())
332+
.build())
333+
.getFeatureSetsList().stream();
334+
} catch (InvalidProtocolBufferException e) {
335+
throw new RuntimeException(
336+
String.format(
337+
"Couldn't fetch featureSets for subscription %s. Reason: %s",
338+
subscription, e.getMessage()));
339+
}
340+
})
332341
.distinct()
333-
.map(wrapException(FeatureSet::toProto))
334342
.collect(Collectors.toList());
335343
}
336344

337345
@Scheduled(fixedDelayString = "${feast.stream.specsOptions.notifyIntervalMilliseconds}")
338-
public void notifyJobsWhenFeatureSetUpdated() {
346+
public void notifyJobsWhenFeatureSetUpdated() throws InvalidProtocolBufferException {
339347
List<FeatureSet> pendingFeatureSets =
340-
featureSetRepository.findAllByStatus(FeatureSetProto.FeatureSetStatus.STATUS_PENDING);
348+
specService
349+
.listFeatureSets(
350+
CoreServiceProto.ListFeatureSetsRequest.Filter.newBuilder()
351+
.setProject("*")
352+
.setFeatureSetName("*")
353+
.setStatus(FeatureSetProto.FeatureSetStatus.STATUS_PENDING)
354+
.build())
355+
.getFeatureSetsList();
341356

342357
pendingFeatureSets.stream()
343358
.map(this::allocateFeatureSetToJobs)
344359
.map(
345360
fs -> {
346361
FeatureSetReference ref =
347-
FeatureSetReference.of(fs.getProject().getName(), fs.getName());
362+
FeatureSetReference.of(fs.getSpec().getProject(), fs.getSpec().getName());
348363
List<FeatureSetDeliveryStatus> deliveryStatuses =
349364
jobRepository.findByFeatureSetReference(ref).stream()
350365
.filter(Job::isRunning)
@@ -360,13 +375,17 @@ public void notifyJobsWhenFeatureSetUpdated() {
360375
&& pair.getRight().stream()
361376
.anyMatch(
362377
jobStatus ->
363-
jobStatus.getDeliveredVersion() < pair.getLeft().getVersion()))
378+
jobStatus.getDeliveredVersion()
379+
< pair.getLeft().getSpec().getVersion()))
364380
.forEach(
365381
pair -> {
366382
FeatureSet fs = pair.getLeft();
367383
List<FeatureSetDeliveryStatus> deliveryStatuses = pair.getRight();
368384

369-
log.info("Sending new FeatureSet {} to Ingestion", fs.getReference());
385+
FeatureSetReference ref =
386+
FeatureSetReference.of(fs.getSpec().getProject(), fs.getSpec().getName());
387+
388+
log.info("Sending new FeatureSet {} to Ingestion", ref);
370389

371390
// Sending latest version of FeatureSet to all currently running IngestionJobs
372391
// (there's one topic for all sets).
@@ -375,7 +394,7 @@ public void notifyJobsWhenFeatureSetUpdated() {
375394
// again later.
376395
try {
377396
specPublisher
378-
.sendDefault(fs.getReference(), fs.toProto().getSpec())
397+
.sendDefault(ref.getReference(), fs.getSpec())
379398
.get(SPEC_PUBLISHING_TIMEOUT_SECONDS, TimeUnit.SECONDS);
380399
} catch (Exception e) {
381400
log.error(
@@ -393,7 +412,7 @@ public void notifyJobsWhenFeatureSetUpdated() {
393412
jobStatus -> {
394413
jobStatus.setDeliveryStatus(
395414
FeatureSetProto.FeatureSetJobDeliveryStatus.STATUS_IN_PROGRESS);
396-
jobStatus.setDeliveredVersion(fs.getVersion());
415+
jobStatus.setDeliveredVersion(fs.getSpec().getVersion());
397416
});
398417
});
399418
}
@@ -412,13 +431,19 @@ public void notifyJobsWhenFeatureSetUpdated() {
412431
@KafkaListener(
413432
topics = {"${feast.stream.specsOptions.specsAckTopic}"},
414433
containerFactory = "kafkaAckListenerContainerFactory")
415-
public void listenAckFromJobs(
416-
ConsumerRecord<String, IngestionJobProto.FeatureSetSpecAck> record) {
434+
public void listenAckFromJobs(ConsumerRecord<String, IngestionJobProto.FeatureSetSpecAck> record)
435+
throws InvalidProtocolBufferException {
417436
String setReference = record.key();
418437
Pair<String, String> projectAndSetName = parseReference(setReference);
419438
FeatureSet featureSet =
420-
featureSetRepository.findFeatureSetByNameAndProject_Name(
421-
projectAndSetName.getRight(), projectAndSetName.getLeft());
439+
specService
440+
.getFeatureSet(
441+
CoreServiceProto.GetFeatureSetRequest.newBuilder()
442+
.setProject(projectAndSetName.getLeft())
443+
.setName(projectAndSetName.getRight())
444+
.build())
445+
.getFeatureSet();
446+
422447
if (featureSet == null) {
423448
log.warn(
424449
String.format("ACKListener received message for unknown FeatureSet %s", setReference));
@@ -427,18 +452,18 @@ public void listenAckFromJobs(
427452

428453
int ackVersion = record.value().getFeatureSetVersion();
429454

430-
if (featureSet.getVersion() != ackVersion) {
455+
if (featureSet.getSpec().getVersion() != ackVersion) {
431456
log.warn(
432457
String.format(
433458
"ACKListener received outdated ack for %s. Current %d, Received %d",
434-
setReference, featureSet.getVersion(), ackVersion));
459+
setReference, featureSet.getSpec().getVersion(), ackVersion));
435460
return;
436461
}
437462

438-
log.info("Updating featureSet {} delivery statuses.", featureSet.getReference());
439-
440463
FeatureSetReference ref =
441-
FeatureSetReference.of(featureSet.getProject().getName(), featureSet.getName());
464+
FeatureSetReference.of(featureSet.getSpec().getProject(), featureSet.getSpec().getName());
465+
466+
log.info("Updating featureSet {} delivery statuses.", ref);
442467

443468
jobRepository
444469
.findById(record.value().getJobName())
@@ -459,10 +484,17 @@ public void listenAckFromJobs(
459484
.equals(FeatureSetProto.FeatureSetJobDeliveryStatus.STATUS_DELIVERED));
460485

461486
if (allDelivered) {
462-
log.info("FeatureSet {} update is completely delivered", featureSet.getReference());
463-
464-
featureSet.setStatus(FeatureSetProto.FeatureSetStatus.STATUS_READY);
465-
featureSetRepository.saveAndFlush(featureSet);
487+
log.info("FeatureSet {} update is completely delivered", ref);
488+
489+
specService.updateFeatureSetStatus(
490+
CoreServiceProto.UpdateFeatureSetStatusRequest.newBuilder()
491+
.setReference(
492+
FeatureSetReferenceProto.FeatureSetReference.newBuilder()
493+
.setName(ref.getFeatureSetName())
494+
.setProject(ref.getProjectName())
495+
.build())
496+
.setStatus(FeatureSetProto.FeatureSetStatus.STATUS_READY)
497+
.build());
466498
}
467499
}
468500
}

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

Lines changed: 31 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,8 @@
3737
import feast.proto.core.CoreServiceProto.ListStoresRequest;
3838
import feast.proto.core.CoreServiceProto.ListStoresResponse;
3939
import feast.proto.core.CoreServiceProto.ListStoresResponse.Builder;
40+
import feast.proto.core.CoreServiceProto.UpdateFeatureSetStatusRequest;
41+
import feast.proto.core.CoreServiceProto.UpdateFeatureSetStatusResponse;
4042
import feast.proto.core.CoreServiceProto.UpdateStoreRequest;
4143
import feast.proto.core.CoreServiceProto.UpdateStoreResponse;
4244
import feast.proto.core.FeatureSetProto;
@@ -89,29 +91,33 @@ public SpecService(
8991
*/
9092
public GetFeatureSetResponse getFeatureSet(GetFeatureSetRequest request)
9193
throws InvalidProtocolBufferException {
94+
FeatureSet featureSet = getFeatureSet(request.getProject(), request.getName());
9295

96+
return GetFeatureSetResponse.newBuilder().setFeatureSet(featureSet.toProto()).build();
97+
}
98+
99+
private FeatureSet getFeatureSet(String projectName, String featureSetName) {
93100
// Validate input arguments
94-
checkValidCharacters(request.getName(), "featureSetName");
101+
checkValidCharacters(featureSetName, "featureSetName");
95102

96-
if (request.getName().isEmpty()) {
103+
if (featureSetName.isEmpty()) {
97104
throw new IllegalArgumentException("No feature set name provided");
98105
}
99106
// Autofill default project if project is not specified
100-
if (request.getProject().isEmpty()) {
101-
request = request.toBuilder().setProject(Project.DEFAULT_NAME).build();
107+
if (projectName.isEmpty()) {
108+
projectName = Project.DEFAULT_NAME;
102109
}
103110

104111
FeatureSet featureSet;
105112

106113
featureSet =
107-
featureSetRepository.findFeatureSetByNameAndProject_Name(
108-
request.getName(), request.getProject());
114+
featureSetRepository.findFeatureSetByNameAndProject_Name(featureSetName, projectName);
109115

110116
if (featureSet == null) {
111117
throw new RetrievalException(
112-
String.format("Feature set with name \"%s\" could not be found.", request.getName()));
118+
String.format("Feature set with name \"%s\" could not be found.", featureSetName));
113119
}
114-
return GetFeatureSetResponse.newBuilder().setFeatureSet(featureSet.toProto()).build();
120+
return featureSet;
115121
}
116122

117123
/**
@@ -138,6 +144,7 @@ public ListFeatureSetsResponse listFeatureSets(ListFeatureSetsRequest.Filter fil
138144
String name = filter.getFeatureSetName();
139145
String project = filter.getProject();
140146
Map<String, String> labelsFilter = filter.getLabelsMap();
147+
FeatureSetStatus statusFilter = filter.getStatus();
141148

142149
if (name.isEmpty()) {
143150
throw new IllegalArgumentException(
@@ -195,6 +202,10 @@ public ListFeatureSetsResponse listFeatureSets(ListFeatureSetsRequest.Filter fil
195202
if (featureSets.size() > 0) {
196203
featureSets =
197204
featureSets.stream()
205+
.filter(
206+
featureSet ->
207+
statusFilter.equals(FeatureSetStatus.STATUS_INVALID)
208+
|| featureSet.getStatus().equals(statusFilter))
198209
.filter(featureSet -> featureSet.hasAllLabels(labelsFilter))
199210
.collect(Collectors.toList());
200211
for (FeatureSet featureSet : featureSets) {
@@ -264,6 +275,18 @@ public ListFeaturesResponse listFeatures(ListFeaturesRequest.Filter filter) {
264275
}
265276
}
266277

278+
/** Update FeatureSet's status by given FeatureSetReference and new status */
279+
public UpdateFeatureSetStatusResponse updateFeatureSetStatus(
280+
UpdateFeatureSetStatusRequest request) {
281+
FeatureSet featureSet =
282+
getFeatureSet(request.getReference().getProject(), request.getReference().getName());
283+
284+
featureSet.setStatus(request.getStatus());
285+
featureSetRepository.saveAndFlush(featureSet);
286+
287+
return UpdateFeatureSetStatusResponse.newBuilder().build();
288+
}
289+
267290
/**
268291
* Get stores matching the store name provided in the filter. If the store name is not provided,
269292
* the method will return all stores currently registered to Feast.

core/src/test/java/feast/core/auth/CoreServiceAuthenticationIT.java

Lines changed: 1 addition & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -162,12 +162,7 @@ void canApplyFeatureSetIfAuthenticated() {
162162
SimpleAPIClient secureApiClient =
163163
getSecureApiClient("AuthenticatedUserWithoutAuthorization@example.com");
164164
FeatureSetProto.FeatureSet expectedFeatureSet =
165-
DataGenerator.createFeatureSet(
166-
DataGenerator.getDefaultSource(),
167-
"project_1",
168-
"test_1",
169-
Collections.emptyList(),
170-
Collections.emptyList());
165+
DataGenerator.createFeatureSet(DataGenerator.getDefaultSource(), "project_1", "test_1");
171166

172167
secureApiClient.simpleApplyFeatureSet(expectedFeatureSet);
173168

0 commit comments

Comments
 (0)