Skip to content

Commit 04d2b47

Browse files
authored
Apply grpc tracing interceptor on online serving (#1242)
Signed-off-by: Khor Shu Heng <khor.heng@gojek.com> Co-authored-by: Khor Shu Heng <khor.heng@gojek.com>
1 parent 17edb99 commit 04d2b47

4 files changed

Lines changed: 127 additions & 103 deletions

File tree

serving/pom.xml

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -155,21 +155,26 @@
155155
<groupId>joda-time</groupId>
156156
<artifactId>joda-time</artifactId>
157157
</dependency>
158-
<!--compile 'io.jaegertracing:jaeger-client:0.31.0'-->
158+
<!--compile 'io.jaegertracing:jaeger-client:1.3.2'-->
159159
<dependency>
160160
<groupId>io.jaegertracing</groupId>
161161
<artifactId>jaeger-client</artifactId>
162-
<version>0.31.0</version>
162+
<version>1.3.2</version>
163163
</dependency>
164164
<dependency>
165165
<groupId>io.opentracing</groupId>
166166
<artifactId>opentracing-api</artifactId>
167-
<version>0.31.0</version>
167+
<version>0.33.0</version>
168168
</dependency>
169169
<dependency>
170170
<groupId>io.opentracing</groupId>
171171
<artifactId>opentracing-noop</artifactId>
172-
<version>0.31.0</version>
172+
<version>0.33.0</version>
173+
</dependency>
174+
<dependency>
175+
<groupId>io.opentracing.contrib</groupId>
176+
<artifactId>opentracing-grpc</artifactId>
177+
<version>0.2.3</version>
173178
</dependency>
174179

175180
<!-- The client -->

serving/src/main/java/feast/serving/config/InstrumentationConfig.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717
package feast.serving.config;
1818

1919
import io.opentracing.Tracer;
20+
import io.opentracing.contrib.grpc.TracingServerInterceptor;
2021
import io.opentracing.noop.NoopTracerFactory;
2122
import io.prometheus.client.exporter.MetricsServlet;
2223
import io.prometheus.client.hotspot.DefaultExports;
@@ -54,4 +55,9 @@ public Tracer tracer() {
5455
return io.jaegertracing.Configuration.fromEnv(feastProperties.getTracing().getServiceName())
5556
.getTracer();
5657
}
58+
59+
@Bean
60+
public TracingServerInterceptor tracingInterceptor(Tracer tracer) {
61+
return TracingServerInterceptor.newBuilder().withTracer(tracer).build();
62+
}
5763
}

serving/src/main/java/feast/serving/controller/ServingServiceGRpcController.java

Lines changed: 12 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -30,16 +30,21 @@
3030
import feast.serving.util.RequestHelper;
3131
import io.grpc.Status;
3232
import io.grpc.stub.StreamObserver;
33-
import io.opentracing.Scope;
3433
import io.opentracing.Span;
3534
import io.opentracing.Tracer;
35+
import io.opentracing.contrib.grpc.TracingServerInterceptor;
3636
import net.devh.boot.grpc.server.service.GrpcService;
3737
import org.slf4j.Logger;
3838
import org.springframework.beans.factory.annotation.Autowired;
3939
import org.springframework.security.access.AccessDeniedException;
4040
import org.springframework.security.core.context.SecurityContextHolder;
4141

42-
@GrpcService(interceptors = {GrpcMessageInterceptor.class, GrpcMonitoringInterceptor.class})
42+
@GrpcService(
43+
interceptors = {
44+
TracingServerInterceptor.class,
45+
GrpcMessageInterceptor.class,
46+
GrpcMonitoringInterceptor.class
47+
})
4348
public class ServingServiceGRpcController extends ServingServiceImplBase {
4449

4550
private static final Logger log =
@@ -75,16 +80,19 @@ public void getFeastServingInfo(
7580
public void getOnlineFeaturesV2(
7681
ServingAPIProto.GetOnlineFeaturesRequestV2 request,
7782
StreamObserver<GetOnlineFeaturesResponse> responseObserver) {
78-
Span span = tracer.buildSpan("getOnlineFeaturesV2").start();
79-
try (Scope scope = tracer.scopeManager().activate(span, false)) {
83+
try {
8084
// authorize for the project in request object.
8185
if (request.getProject() != null && !request.getProject().isEmpty()) {
8286
// project set at root level overrides the project set at feature table level
8387
this.authorizationService.authorizeRequest(
8488
SecurityContextHolder.getContext(), request.getProject());
8589
}
8690
RequestHelper.validateOnlineRequest(request);
91+
Span span = tracer.buildSpan("getOnlineFeaturesV2").start();
8792
GetOnlineFeaturesResponse onlineFeatures = servingServiceV2.getOnlineFeatures(request);
93+
if (span != null) {
94+
span.finish();
95+
}
8896
responseObserver.onNext(onlineFeatures);
8997
responseObserver.onCompleted();
9098
} catch (SpecRetrievalException e) {
@@ -102,6 +110,5 @@ public void getOnlineFeaturesV2(
102110
log.warn("Failed to get Online Features", e);
103111
responseObserver.onError(e);
104112
}
105-
span.finish();
106113
}
107114
}

serving/src/main/java/feast/serving/service/OnlineServingServiceV2.java

Lines changed: 100 additions & 94 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,7 @@
3232
import feast.storage.api.retriever.Feature;
3333
import feast.storage.api.retriever.OnlineRetrieverV2;
3434
import io.grpc.Status;
35-
import io.opentracing.Scope;
35+
import io.opentracing.Span;
3636
import io.opentracing.Tracer;
3737
import java.util.*;
3838
import java.util.stream.Collectors;
@@ -71,105 +71,111 @@ public GetOnlineFeaturesResponse getOnlineFeatures(GetOnlineFeaturesRequestV2 re
7171
projectName = "default";
7272
}
7373

74-
try (Scope scope = tracer.buildSpan("getOnlineFeaturesV2").startActive(true)) {
75-
List<GetOnlineFeaturesRequestV2.EntityRow> entityRows = request.getEntityRowsList();
76-
// Collect the feature/entity value for each entity row in entityValueMap
77-
Map<GetOnlineFeaturesRequestV2.EntityRow, Map<String, ValueProto.Value>> entityValuesMap =
78-
entityRows.stream().collect(Collectors.toMap(row -> row, row -> new HashMap<>()));
79-
// Collect the feature/entity status metadata for each entity row in entityValueMap
80-
Map<GetOnlineFeaturesRequestV2.EntityRow, Map<String, GetOnlineFeaturesResponse.FieldStatus>>
81-
entityStatusesMap =
82-
entityRows.stream().collect(Collectors.toMap(row -> row, row -> new HashMap<>()));
83-
84-
entityRows.forEach(
85-
entityRow -> {
86-
Map<String, ValueProto.Value> valueMap = entityRow.getFieldsMap();
87-
entityValuesMap.get(entityRow).putAll(valueMap);
88-
entityStatusesMap.get(entityRow).putAll(getMetadataMap(valueMap, false, false));
89-
});
90-
91-
List<List<Optional<Feature>>> entityRowsFeatures =
92-
retriever.getOnlineFeatures(projectName, entityRows, featureReferences);
93-
94-
if (entityRowsFeatures.size() != entityRows.size()) {
95-
throw Status.INTERNAL
96-
.withDescription(
97-
"The no. of FeatureRow obtained from OnlineRetriever"
98-
+ "does not match no. of entityRow passed.")
99-
.asRuntimeException();
100-
}
74+
List<GetOnlineFeaturesRequestV2.EntityRow> entityRows = request.getEntityRowsList();
75+
// Collect the feature/entity value for each entity row in entityValueMap
76+
Map<GetOnlineFeaturesRequestV2.EntityRow, Map<String, ValueProto.Value>> entityValuesMap =
77+
entityRows.stream().collect(Collectors.toMap(row -> row, row -> new HashMap<>()));
78+
// Collect the feature/entity status metadata for each entity row in entityValueMap
79+
Map<GetOnlineFeaturesRequestV2.EntityRow, Map<String, GetOnlineFeaturesResponse.FieldStatus>>
80+
entityStatusesMap =
81+
entityRows.stream().collect(Collectors.toMap(row -> row, row -> new HashMap<>()));
82+
83+
entityRows.forEach(
84+
entityRow -> {
85+
Map<String, ValueProto.Value> valueMap = entityRow.getFieldsMap();
86+
entityValuesMap.get(entityRow).putAll(valueMap);
87+
entityStatusesMap.get(entityRow).putAll(getMetadataMap(valueMap, false, false));
88+
});
10189

102-
for (int i = 0; i < entityRows.size(); i++) {
103-
GetOnlineFeaturesRequestV2.EntityRow entityRow = entityRows.get(i);
104-
List<Optional<Feature>> curEntityRowFeatures = entityRowsFeatures.get(i);
105-
106-
Map<FeatureReferenceV2, Optional<Feature>> featureReferenceFeatureMap =
107-
getFeatureRefFeatureMap(curEntityRowFeatures);
108-
109-
Map<String, ValueProto.Value> allValueMaps = new HashMap<>();
110-
Map<String, GetOnlineFeaturesResponse.FieldStatus> allStatusMaps = new HashMap<>();
111-
112-
for (FeatureReferenceV2 featureReference : featureReferences) {
113-
if (featureReferenceFeatureMap.containsKey(featureReference)) {
114-
Optional<Feature> feature = featureReferenceFeatureMap.get(featureReference);
115-
116-
FeatureTableSpec featureTableSpec =
117-
specService.getFeatureTableSpec(projectName, feature.get().getFeatureReference());
118-
FeatureProto.FeatureSpecV2 featureSpec =
119-
specService.getFeatureSpec(projectName, feature.get().getFeatureReference());
120-
ValueProto.ValueType.Enum valueTypeEnum = featureSpec.getValueType();
121-
ValueProto.Value.ValCase valueCase = feature.get().getFeatureValue().getValCase();
122-
boolean isMatchingFeatureSpec = checkSameFeatureSpec(valueTypeEnum, valueCase);
123-
124-
boolean isOutsideMaxAge = checkOutsideMaxAge(featureTableSpec, entityRow, feature);
125-
Map<String, ValueProto.Value> valueMap =
126-
unpackValueMap(feature, isOutsideMaxAge, isMatchingFeatureSpec);
127-
allValueMaps.putAll(valueMap);
128-
129-
// Generate metadata for feature values and merge into entityFieldsMap
130-
Map<String, GetOnlineFeaturesResponse.FieldStatus> statusMap =
131-
getMetadataMap(valueMap, !isMatchingFeatureSpec, isOutsideMaxAge);
132-
allStatusMaps.putAll(statusMap);
133-
134-
// Populate metrics/log request
135-
populateCountMetrics(statusMap, projectName);
136-
} else {
137-
Map<String, ValueProto.Value> valueMap =
138-
new HashMap<>() {
139-
{
140-
put(
141-
FeatureV2.getFeatureStringRef(featureReference),
142-
ValueProto.Value.newBuilder().build());
143-
}
144-
};
145-
allValueMaps.putAll(valueMap);
90+
Span onlineRetrievalSpan = tracer.buildSpan("onlineRetrieval").start();
91+
if (onlineRetrievalSpan != null) {
92+
onlineRetrievalSpan.setTag("entities", entityRows.size());
93+
onlineRetrievalSpan.setTag("features", featureReferences.size());
94+
}
95+
List<List<Optional<Feature>>> entityRowsFeatures =
96+
retriever.getOnlineFeatures(projectName, entityRows, featureReferences);
97+
if (onlineRetrievalSpan != null) {
98+
onlineRetrievalSpan.finish();
99+
}
146100

147-
Map<String, GetOnlineFeaturesResponse.FieldStatus> statusMap =
148-
getMetadataMap(valueMap, true, false);
149-
allStatusMaps.putAll(statusMap);
101+
if (entityRowsFeatures.size() != entityRows.size()) {
102+
throw Status.INTERNAL
103+
.withDescription(
104+
"The no. of FeatureRow obtained from OnlineRetriever"
105+
+ "does not match no. of entityRow passed.")
106+
.asRuntimeException();
107+
}
150108

151-
// Populate metrics/log request
152-
populateCountMetrics(statusMap, projectName);
153-
}
109+
for (int i = 0; i < entityRows.size(); i++) {
110+
GetOnlineFeaturesRequestV2.EntityRow entityRow = entityRows.get(i);
111+
List<Optional<Feature>> curEntityRowFeatures = entityRowsFeatures.get(i);
112+
113+
Map<FeatureReferenceV2, Optional<Feature>> featureReferenceFeatureMap =
114+
getFeatureRefFeatureMap(curEntityRowFeatures);
115+
116+
Map<String, ValueProto.Value> allValueMaps = new HashMap<>();
117+
Map<String, GetOnlineFeaturesResponse.FieldStatus> allStatusMaps = new HashMap<>();
118+
119+
for (FeatureReferenceV2 featureReference : featureReferences) {
120+
if (featureReferenceFeatureMap.containsKey(featureReference)) {
121+
Optional<Feature> feature = featureReferenceFeatureMap.get(featureReference);
122+
123+
FeatureTableSpec featureTableSpec =
124+
specService.getFeatureTableSpec(projectName, feature.get().getFeatureReference());
125+
FeatureProto.FeatureSpecV2 featureSpec =
126+
specService.getFeatureSpec(projectName, feature.get().getFeatureReference());
127+
ValueProto.ValueType.Enum valueTypeEnum = featureSpec.getValueType();
128+
ValueProto.Value.ValCase valueCase = feature.get().getFeatureValue().getValCase();
129+
boolean isMatchingFeatureSpec = checkSameFeatureSpec(valueTypeEnum, valueCase);
130+
131+
boolean isOutsideMaxAge = checkOutsideMaxAge(featureTableSpec, entityRow, feature);
132+
Map<String, ValueProto.Value> valueMap =
133+
unpackValueMap(feature, isOutsideMaxAge, isMatchingFeatureSpec);
134+
allValueMaps.putAll(valueMap);
135+
136+
// Generate metadata for feature values and merge into entityFieldsMap
137+
Map<String, GetOnlineFeaturesResponse.FieldStatus> statusMap =
138+
getMetadataMap(valueMap, !isMatchingFeatureSpec, isOutsideMaxAge);
139+
allStatusMaps.putAll(statusMap);
140+
141+
// Populate metrics/log request
142+
populateCountMetrics(statusMap, projectName);
143+
} else {
144+
Map<String, ValueProto.Value> valueMap =
145+
new HashMap<>() {
146+
{
147+
put(
148+
FeatureV2.getFeatureStringRef(featureReference),
149+
ValueProto.Value.newBuilder().build());
150+
}
151+
};
152+
allValueMaps.putAll(valueMap);
153+
154+
Map<String, GetOnlineFeaturesResponse.FieldStatus> statusMap =
155+
getMetadataMap(valueMap, true, false);
156+
allStatusMaps.putAll(statusMap);
157+
158+
// Populate metrics/log request
159+
populateCountMetrics(statusMap, projectName);
154160
}
155-
entityValuesMap.get(entityRow).putAll(allValueMaps);
156-
entityStatusesMap.get(entityRow).putAll(allStatusMaps);
157161
}
158-
159-
// Build response field values from entityValuesMap and entityStatusesMap
160-
// Response field values should be in the same order as the entityRows provided by the user.
161-
List<GetOnlineFeaturesResponse.FieldValues> fieldValuesList =
162-
entityRows.stream()
163-
.map(
164-
entityRow -> {
165-
return GetOnlineFeaturesResponse.FieldValues.newBuilder()
166-
.putAllFields(entityValuesMap.get(entityRow))
167-
.putAllStatuses(entityStatusesMap.get(entityRow))
168-
.build();
169-
})
170-
.collect(Collectors.toList());
171-
return GetOnlineFeaturesResponse.newBuilder().addAllFieldValues(fieldValuesList).build();
162+
entityValuesMap.get(entityRow).putAll(allValueMaps);
163+
entityStatusesMap.get(entityRow).putAll(allStatusMaps);
172164
}
165+
166+
// Build response field values from entityValuesMap and entityStatusesMap
167+
// Response field values should be in the same order as the entityRows provided by the user.
168+
List<GetOnlineFeaturesResponse.FieldValues> fieldValuesList =
169+
entityRows.stream()
170+
.map(
171+
entityRow -> {
172+
return GetOnlineFeaturesResponse.FieldValues.newBuilder()
173+
.putAllFields(entityValuesMap.get(entityRow))
174+
.putAllStatuses(entityStatusesMap.get(entityRow))
175+
.build();
176+
})
177+
.collect(Collectors.toList());
178+
return GetOnlineFeaturesResponse.newBuilder().addAllFieldValues(fieldValuesList).build();
173179
}
174180

175181
private boolean checkSameFeatureSpec(

0 commit comments

Comments
 (0)