1616
1717package feast .serving .service ;
1818
19+ import static feast .serving .util .Metrics .missingKeyCount ;
20+ import static feast .serving .util .Metrics .requestLatency ;
21+ import static feast .serving .util .Metrics .requestCount ;
22+ import static feast .serving .util .Metrics .staleKeyCount ;
23+
1924import com .google .common .collect .Maps ;
2025import com .google .protobuf .AbstractMessageLite ;
2126import com .google .protobuf .Duration ;
4146import io .grpc .Status ;
4247import io .opentracing .Scope ;
4348import io .opentracing .Tracer ;
49+ import io .prometheus .client .Histogram .Timer ;
4450import java .util .List ;
4551import java .util .Map ;
4652import java .util .stream .Collectors ;
@@ -61,7 +67,9 @@ public RedisServingService(JedisPool jedisPool, CachedSpecService specService, T
6167 this .tracer = tracer ;
6268 }
6369
64- /** {@inheritDoc} */
70+ /**
71+ * {@inheritDoc}
72+ */
6573 @ Override
6674 public GetFeastServingInfoResponse getFeastServingInfo (
6775 GetFeastServingInfoRequest getFeastServingInfoRequest ) {
@@ -70,10 +78,13 @@ public GetFeastServingInfoResponse getFeastServingInfo(
7078 .build ();
7179 }
7280
73- /** {@inheritDoc} */
81+ /**
82+ * {@inheritDoc}
83+ */
7484 @ Override
7585 public GetOnlineFeaturesResponse getOnlineFeatures (GetOnlineFeaturesRequest request ) {
7686 try (Scope scope = tracer .buildSpan ("Redis-getOnlineFeatures" ).startActive (true )) {
87+ Timer getOnlineFeaturesTimer = requestLatency .labels ("getOnlineFeatures" ).startTimer ();
7788 GetOnlineFeaturesResponse .Builder getOnlineFeaturesResponseBuilder =
7889 GetOnlineFeaturesResponse .newBuilder ();
7990
@@ -114,6 +125,7 @@ public GetOnlineFeaturesResponse getOnlineFeatures(GetOnlineFeaturesRequest requ
114125 featureValuesMap .values ().stream ()
115126 .map (m -> FieldValues .newBuilder ().putAllFields (m ).build ())
116127 .collect (Collectors .toList ());
128+ getOnlineFeaturesTimer .observeDuration ();
117129 return getOnlineFeaturesResponseBuilder .addAllFieldValues (fieldValues ).build ();
118130 }
119131 }
@@ -166,9 +178,11 @@ private RedisKey makeRedisKey(
166178 for (int i = 0 ; i < featureSetEntityNames .size (); i ++) {
167179 String entityName = featureSetEntityNames .get (i );
168180
169- if (!fieldsMap .containsKey (entityName )){
181+ if (!fieldsMap .containsKey (entityName )) {
170182 throw Status .INVALID_ARGUMENT
171- .withDescription (String .format ("Entity row fields \" %s\" does not contain required entity field \" %s\" " , fieldsMap .keySet ().toString (), entityName ))
183+ .withDescription (String
184+ .format ("Entity row fields \" %s\" does not contain required entity field \" %s\" " ,
185+ fieldsMap .keySet ().toString (), entityName ))
172186 .asRuntimeException ();
173187 }
174188
@@ -186,33 +200,45 @@ private void sendAndProcessMultiGet(
186200 throws InvalidProtocolBufferException {
187201
188202 List <byte []> jedisResps = sendMultiGet (redisKeys );
189-
203+ Timer processResponseTimer = requestLatency .labels ("processResponse" )
204+ .startTimer ();
190205 try (Scope scope = tracer .buildSpan ("Redis-processResponse" ).startActive (true )) {
191206 String featureSetId =
192207 String .format ("%s:%d" , featureSetRequest .getName (), featureSetRequest .getVersion ());
208+
193209 Map <String , Value > nullValues =
194210 featureSetRequest .getFeatureNamesList ().stream ()
195211 .collect (
196212 Collectors .toMap (
197213 name -> featureSetId + ":" + name , name -> Value .newBuilder ().build ()));
214+
198215 for (int i = 0 ; i < jedisResps .size (); i ++) {
199216 EntityRow entityRow = entityRows .get (i );
200217 Map <String , Value > featureValues = featureValuesMap .get (entityRow );
218+
201219 byte [] jedisResponse = jedisResps .get (i );
202220 if (jedisResponse == null ) {
221+ missingKeyCount .labels (featureSetRequest .getName ()).inc ();
203222 featureValues .putAll (nullValues );
204223 continue ;
205224 }
225+
206226 FeatureRow featureRow = FeatureRow .parseFrom (jedisResponse );
227+
207228 boolean stale = isStale (featureSetRequest , entityRow , featureRow );
208229 if (stale ) {
230+ staleKeyCount .labels (featureSetRequest .getName ()).inc ();
209231 featureValues .putAll (nullValues );
210232 continue ;
211233 }
234+
235+ requestCount .labels (featureSetRequest .getName ()).inc ();
212236 featureRow .getFieldsList ().stream ()
213237 .filter (f -> featureSetRequest .getFeatureNamesList ().contains (f .getName ()))
214238 .forEach (f -> featureValues .put (featureSetId + ":" + f .getName (), f .getValue ()));
215239 }
240+ } finally {
241+ processResponseTimer .observeDuration ();
216242 }
217243 }
218244
@@ -237,6 +263,7 @@ private boolean isStale(
237263 */
238264 private List <byte []> sendMultiGet (List <RedisKey > keys ) {
239265 try (Scope scope = tracer .buildSpan ("Redis-sendMultiGet" ).startActive (true )) {
266+ Timer sendMultiGetTimer = requestLatency .labels ("sendMultiGet" ).startTimer ();
240267 try (Jedis jedis = jedisPool .getResource ()) {
241268 byte [][] binaryKeys =
242269 keys .stream ()
@@ -249,6 +276,8 @@ private List<byte[]> sendMultiGet(List<RedisKey> keys) {
249276 .withDescription ("Unable to retrieve feature from Redis" )
250277 .withCause (e )
251278 .asRuntimeException ();
279+ } finally {
280+ sendMultiGetTimer .observeDuration ();
252281 }
253282 }
254283 }
0 commit comments