1+ package feast
2+
3+ import (
4+ "context"
5+ "github.com/feast-dev/feast/go/protos/feast/types"
6+ "github.com/feast-dev/feast/go/protos/feast/third_party/grpc/connector"
7+ )
8+
9+ // GRPCClient is an implementation of KV that talks over RPC.
10+ type GRPCClient struct { client connector.OnlineStoreClient
11+ destructor func () }
12+
13+ func (m * GRPCClient ) OnlineRead (entityKeys []types.EntityKey , view string , features []string ) ([][]Feature , error ) {
14+ entityKeysRef := make ([]* types.EntityKey , len (entityKeys ))
15+ for i := 0 ; i < len (entityKeys ); i ++ {
16+ entityKeysRef [i ] = & entityKeys [i ]
17+ }
18+ results , err := m .client .OnlineRead (context .Background (), & connector.OnlineReadRequest {
19+ EntityKeys : entityKeysRef ,
20+ View : view ,
21+ Features : features ,
22+ })
23+ if err != nil {
24+ return nil , err
25+ }
26+ feature2D := results .GetResults ()
27+ featureResults := make ([][]Feature , len (feature2D ))
28+ for entityIndex , featureList := range feature2D {
29+ connectorList := featureList .GetFeatureList ()
30+ featureResults [entityIndex ] = make ([]Feature , len (connectorList ))
31+ for featureIndex , feature := range connectorList {
32+ featureResults [entityIndex ][featureIndex ] = Feature { reference : * feature .GetReference (),
33+ timestamp : * feature .GetTimestamp (),
34+ value : * feature .GetValue () }
35+ }
36+ }
37+ return featureResults , nil
38+ }
39+
40+ func (m * GRPCClient ) Destruct () {
41+ m .destructor ()
42+ }
43+
44+ // Here is the gRPC server that GRPCClient talks to.
45+ type GRPCServer struct {
46+ // This is the real implementation
47+ Impl OnlineStore
48+ connector.UnimplementedOnlineStoreServer
49+ }
50+
51+ func (m * GRPCServer ) OnlineRead (
52+ ctx context.Context ,
53+ req * connector.OnlineReadRequest ) (* connector.OnlineReadResponse , error ) {
54+ numEntityKeys := len (req .EntityKeys )
55+ entityKeys := make ([]types.EntityKey , numEntityKeys )
56+ for i := 0 ; i < numEntityKeys ; i ++ {
57+ entityKeys [i ] = * req .EntityKeys [i ]
58+ }
59+ features , err := m .Impl .OnlineRead (entityKeys , req .View , req .Features )
60+ if err != nil {
61+ return nil , err
62+ }
63+ response := connector.OnlineReadResponse {Results : make ([]* connector.ConnectorFeatureList , len (features ))}
64+
65+ for entityIndex , featureList := range features {
66+ response .Results [entityIndex ] = & connector.ConnectorFeatureList {FeatureList : make ([]* connector.ConnectorFeature , len (featureList ))}
67+ for featureIndex , feature := range featureList {
68+ reference := feature .reference
69+ value := feature .value
70+ timestamp := feature .timestamp
71+ response .Results [entityIndex ].FeatureList [featureIndex ] = & connector.ConnectorFeature { Reference : & reference ,
72+ Value : & value ,
73+ Timestamp : & timestamp }
74+ }
75+ }
76+ return & response , nil
77+ }
0 commit comments