2929import feast .core .CoreServiceProto .GetStoresRequest ;
3030import feast .core .CoreServiceProto .GetStoresRequest .Filter ;
3131import feast .core .CoreServiceProto .GetStoresResponse ;
32+ import feast .core .CoreServiceProto .UpdateStoreRequest ;
33+ import feast .core .CoreServiceProto .UpdateStoreResponse ;
3234import feast .core .FeatureSetProto .FeatureSetSpec ;
3335import feast .core .StoreProto .Store ;
3436import feast .core .StoreProto .Store .Subscription ;
3537import feast .core .exception .RetrievalException ;
3638import feast .core .service .JobCoordinatorService ;
3739import feast .core .service .SpecService ;
38- import io .grpc .Status ;
39- import io .grpc .StatusRuntimeException ;
4040import io .grpc .stub .StreamObserver ;
4141import java .util .ArrayList ;
42+ import java .util .HashSet ;
4243import java .util .List ;
44+ import java .util .Set ;
4345import java .util .regex .Pattern ;
4446import java .util .stream .Collectors ;
4547import lombok .extern .slf4j .Slf4j ;
4648import org .lognet .springboot .grpc .GRpcService ;
4749import org .springframework .beans .factory .annotation .Autowired ;
4850import org .springframework .transaction .annotation .Transactional ;
4951
50- /** Implementation of the feast core GRPC service. */
52+ /**
53+ * Implementation of the feast core GRPC service.
54+ */
5155@ Slf4j
5256@ GRpcService
5357public class CoreServiceImpl extends CoreServiceImplBase {
5458
55- @ Autowired private SpecService specService ;
56- @ Autowired private StatsDClient statsDClient ;
57- @ Autowired private JobCoordinatorService jobCoordinatorService ;
59+ @ Autowired
60+ private SpecService specService ;
61+ @ Autowired
62+ private StatsDClient statsDClient ;
63+ @ Autowired
64+ private JobCoordinatorService jobCoordinatorService ;
5865
5966 @ Override
6067 public void getFeastCoreVersion (
@@ -72,7 +79,7 @@ public void getFeatureSets(
7279 responseObserver .onNext (response );
7380 responseObserver .onCompleted ();
7481 } catch (RetrievalException | InvalidProtocolBufferException e ) {
75- responseObserver .onError (getRuntimeException ( e ) );
82+ responseObserver .onError (e );
7683 }
7784 }
7885
@@ -85,7 +92,7 @@ public void getStores(
8592 responseObserver .onNext (response );
8693 responseObserver .onCompleted ();
8794 } catch (RetrievalException e ) {
88- responseObserver .onError (getRuntimeException ( e ) );
95+ responseObserver .onError (e );
8996 }
9097 }
9198
@@ -124,17 +131,37 @@ public void applyFeatureSet(
124131 responseObserver .onNext (response );
125132 responseObserver .onCompleted ();
126133 } catch (Exception e ) {
127- responseObserver .onError (getRuntimeException ( e ) );
134+ responseObserver .onError (e );
128135 }
129136 }
130137
131- private StatusRuntimeException getRuntimeException (Exception e ) {
132- return new StatusRuntimeException (
133- Status .fromCode (Status .Code .INTERNAL ).withDescription (e .getMessage ()).withCause (e ));
134- }
138+ @ Override
139+ public void updateStore (UpdateStoreRequest request ,
140+ StreamObserver <UpdateStoreResponse > responseObserver ) {
141+ try {
142+ UpdateStoreResponse response = specService .updateStore (request );
143+ responseObserver .onNext (response );
144+ responseObserver .onCompleted ();
135145
136- private StatusRuntimeException getBadRequestException (Exception e ) {
137- return new StatusRuntimeException (
138- Status .fromCode (Status .Code .OUT_OF_RANGE ).withDescription (e .getMessage ()).withCause (e ));
146+ Set <FeatureSetSpec > featureSetSpecs = new HashSet <>();
147+ Store store = response .getStore ();
148+ for (Subscription subscription : store .getSubscriptionsList ()) {
149+ featureSetSpecs .addAll (
150+ specService .getFeatureSets (
151+ GetFeatureSetsRequest .Filter .newBuilder ()
152+ .setFeatureSetName (subscription .getName ())
153+ .setFeatureSetVersion (subscription .getVersion ())
154+ .build ())
155+ .getFeatureSetsList ()
156+ );
157+ }
158+ featureSetSpecs .stream ()
159+ .collect (Collectors .groupingBy (FeatureSetSpec ::getName ))
160+ .entrySet ()
161+ .stream ()
162+ .forEach (kv -> jobCoordinatorService .startOrUpdateJob (kv .getValue (), store ));
163+ } catch (Exception e ) {
164+ responseObserver .onError (e );
165+ }
139166 }
140167}
0 commit comments