|
32 | 32 | import feast.storage.api.retriever.Feature; |
33 | 33 | import feast.storage.api.retriever.OnlineRetrieverV2; |
34 | 34 | import io.grpc.Status; |
35 | | -import io.opentracing.Scope; |
| 35 | +import io.opentracing.Span; |
36 | 36 | import io.opentracing.Tracer; |
37 | 37 | import java.util.*; |
38 | 38 | import java.util.stream.Collectors; |
@@ -71,105 +71,111 @@ public GetOnlineFeaturesResponse getOnlineFeatures(GetOnlineFeaturesRequestV2 re |
71 | 71 | projectName = "default"; |
72 | 72 | } |
73 | 73 |
|
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 | + }); |
101 | 89 |
|
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 | + } |
146 | 100 |
|
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 | + } |
150 | 108 |
|
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); |
154 | 160 | } |
155 | | - entityValuesMap.get(entityRow).putAll(allValueMaps); |
156 | | - entityStatusesMap.get(entityRow).putAll(allStatusMaps); |
157 | 161 | } |
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); |
172 | 164 | } |
| 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(); |
173 | 179 | } |
174 | 180 |
|
175 | 181 | private boolean checkSameFeatureSpec( |
|
0 commit comments