Skip to content

Commit 4b2f0a8

Browse files
author
Tsotne Tabidze
committed
Add feature logging to http server
Signed-off-by: Tsotne Tabidze <tsotne@tecton.ai>
1 parent 54c3da8 commit 4b2f0a8

4 files changed

Lines changed: 53 additions & 6 deletions

File tree

go/internal/feast/server/grpc_server.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -86,7 +86,7 @@ func (s *grpcServingServiceServer) GetOnlineFeatures(ctx context.Context, reques
8686
fmt.Printf("Couldn't instantiate logger for feature service %s: %+v", featuresOrService.FeatureService.Name, err)
8787
}
8888

89-
err = logger.Log(entityValuesMap, resp.Results[len(request.Entities):], resp.Metadata.FeatureNames.Val[len(request.Entities):], request.RequestContext, requestId)
89+
err = logger.Log(request.Entities, resp.Results[len(request.Entities):], resp.Metadata.FeatureNames.Val[len(request.Entities):], request.RequestContext, requestId)
9090
if err != nil {
9191
fmt.Printf("LoggerImpl error[%s]: %+v", featuresOrService.FeatureService.Name, err)
9292
}

go/internal/feast/server/http_server.go

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,9 @@ import (
77
"github.com/feast-dev/feast/go/internal/feast"
88
"github.com/feast-dev/feast/go/internal/feast/model"
99
"github.com/feast-dev/feast/go/internal/feast/server/logging"
10+
"github.com/feast-dev/feast/go/protos/feast/serving"
1011
prototypes "github.com/feast-dev/feast/go/protos/feast/types"
12+
"github.com/feast-dev/feast/go/types"
1113
"net/http"
1214
)
1315

@@ -152,6 +154,7 @@ func (s *httpServer) getOnlineFeatures(w http.ResponseWriter, r *http.Request) {
152154
featureService, err = s.fs.GetFeatureService(*request.FeatureService)
153155
if err != nil {
154156
http.Error(w, fmt.Sprintf("Error getting feature service from registry: %+v", err), http.StatusInternalServerError)
157+
return
155158
}
156159
}
157160
entitiesProto := make(map[string]*prototypes.RepeatedValue)
@@ -173,6 +176,7 @@ func (s *httpServer) getOnlineFeatures(w http.ResponseWriter, r *http.Request) {
173176

174177
if err != nil {
175178
http.Error(w, fmt.Sprintf("Error getting feature vector: %+v", err), http.StatusInternalServerError)
179+
return
176180
}
177181

178182
var featureNames []string
@@ -209,9 +213,42 @@ func (s *httpServer) getOnlineFeatures(w http.ResponseWriter, r *http.Request) {
209213

210214
if err != nil {
211215
http.Error(w, fmt.Sprintf("Error encoding response: %+v", err), http.StatusInternalServerError)
216+
return
212217
}
213218

214219
w.Header().Set("Content-Type", "application/json")
220+
221+
if featureService != nil && featureService.LoggingConfig != nil && s.loggingService != nil {
222+
logger, err := s.loggingService.GetOrCreateLogger(featureService)
223+
if err != nil {
224+
http.Error(w, fmt.Sprintf("Couldn't instantiate logger for feature service %s: %+v", featureService.Name, err), http.StatusInternalServerError)
225+
return
226+
}
227+
228+
requestId := GenerateRequestId()
229+
230+
// Note: we're converting arrow to proto for feature logging. In the future we should
231+
// base feature logging on arrow so that we don't have to do this extra conversion.
232+
var featureVectorProtos []*serving.GetOnlineFeaturesResponse_FeatureVector
233+
for _, vector := range featureVectors[len(request.Entities):] {
234+
values, err := types.ArrowValuesToProtoValues(vector.Values)
235+
if err != nil {
236+
http.Error(w, fmt.Sprintf("Couldn't convert arrow values into protobuf: %+v", err), http.StatusInternalServerError)
237+
return
238+
}
239+
featureVectorProtos = append(featureVectorProtos, &serving.GetOnlineFeaturesResponse_FeatureVector{
240+
Values: values,
241+
Statuses: vector.Statuses,
242+
EventTimestamps: vector.Timestamps,
243+
})
244+
}
245+
246+
err = logger.Log(entitiesProto, featureVectorProtos, featureNames, requestContextProto, requestId)
247+
if err != nil {
248+
http.Error(w, fmt.Sprintf("LoggerImpl error[%s]: %+v", featureService.Name, err), http.StatusInternalServerError)
249+
return
250+
}
251+
}
215252
}
216253

217254
func (s *httpServer) Serve(host string, port int) error {

go/internal/feast/server/logging/logger.go

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -42,7 +42,7 @@ type LogSink interface {
4242
}
4343

4444
type Logger interface {
45-
Log(joinKeyToEntityValues map[string][]*types.Value, featureVectors []*serving.GetOnlineFeaturesResponse_FeatureVector, featureNames []string, requestData map[string]*types.RepeatedValue, requestId string) error
45+
Log(joinKeyToEntityValues map[string]*types.RepeatedValue, featureVectors []*serving.GetOnlineFeaturesResponse_FeatureVector, featureNames []string, requestData map[string]*types.RepeatedValue, requestId string) error
4646
}
4747

4848
type LoggerImpl struct {
@@ -207,7 +207,7 @@ func getFullFeatureName(featureViewName string, featureName string) string {
207207
return fmt.Sprintf("%s__%s", featureViewName, featureName)
208208
}
209209

210-
func (l *LoggerImpl) Log(joinKeyToEntityValues map[string][]*types.Value, featureVectors []*serving.GetOnlineFeaturesResponse_FeatureVector, featureNames []string, requestData map[string]*types.RepeatedValue, requestId string) error {
210+
func (l *LoggerImpl) Log(joinKeyToEntityValues map[string]*types.RepeatedValue, featureVectors []*serving.GetOnlineFeaturesResponse_FeatureVector, featureNames []string, requestData map[string]*types.RepeatedValue, requestId string) error {
211211
if len(featureVectors) == 0 {
212212
return nil
213213
}
@@ -250,7 +250,7 @@ func (l *LoggerImpl) Log(joinKeyToEntityValues map[string][]*types.Value, featur
250250
if !ok {
251251
return errors.Errorf("Missing join key %s in log data", joinKey)
252252
}
253-
entityValues[idx] = rows[rowIdx]
253+
entityValues[idx] = rows.Val[rowIdx]
254254
}
255255

256256
requestDataValues := make([]*types.Value, len(l.schema.RequestData))
@@ -283,6 +283,6 @@ func (l *LoggerImpl) Log(joinKeyToEntityValues map[string][]*types.Value, featur
283283

284284
type DummyLoggerImpl struct{}
285285

286-
func (l *DummyLoggerImpl) Log(joinKeyToEntityValues map[string][]*types.Value, featureVectors []*serving.GetOnlineFeaturesResponse_FeatureVector, featureNames []string, requestData map[string]*types.RepeatedValue, requestId string) error {
286+
func (l *DummyLoggerImpl) Log(joinKeyToEntityValues map[string]*types.RepeatedValue, featureVectors []*serving.GetOnlineFeaturesResponse_FeatureVector, featureNames []string, requestData map[string]*types.RepeatedValue, requestId string) error {
287287
return nil
288288
}

go/internal/feast/server/logging/logger_test.go

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -90,7 +90,17 @@ func TestLogAndFlushToFile(t *testing.T) {
9090
assert.Nil(t, err)
9191

9292
assert.Nil(t, logger.Log(
93-
map[string][]*types.Value{"driver_id": {{Val: &types.Value_Int32Val{Int32Val: 111}}}},
93+
map[string]*types.RepeatedValue{
94+
"driver_id": {
95+
Val: []*types.Value{
96+
{
97+
Val: &types.Value_Int32Val{
98+
Int32Val: 111,
99+
},
100+
},
101+
},
102+
},
103+
},
94104
[]*serving.GetOnlineFeaturesResponse_FeatureVector{
95105
{
96106
Values: []*types.Value{{Val: &types.Value_DoubleVal{DoubleVal: 2.0}}},

0 commit comments

Comments
 (0)