Skip to content

Commit 472f85b

Browse files
committed
Merge branch '0.3-dev' into 0.3-dev-davidheryanto
2 parents afdc01a + bccdbf7 commit 472f85b

9 files changed

Lines changed: 52 additions & 36 deletions

File tree

ingestion/src/main/java/feast/store/serving/bigquery/FeatureRowExtendedToTableRowDoFn.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22

33
import com.google.api.services.bigquery.model.TableRow;
44
import com.google.protobuf.util.Timestamps;
5-
import feast.types.FeatureProto.Field;
5+
import feast.types.FieldProto.Field;
66
import feast.types.FeatureRowExtendedProto.FeatureRowExtended;
77
import feast.types.FeatureRowProto.FeatureRow;
88
import java.util.Base64;

ingestion/src/main/java/feast/store/serving/redis/FeatureRowToRedisMutationDoFn.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@
2323
import feast.storage.RedisProto.RedisKey.Builder;
2424
import feast.store.serving.redis.RedisCustomIO.Method;
2525
import feast.store.serving.redis.RedisCustomIO.RedisMutation;
26-
import feast.types.FeatureProto.Field;
26+
import feast.types.FieldProto.Field;
2727
import feast.types.FeatureRowExtendedProto.FeatureRowExtended;
2828
import feast.types.FeatureRowProto.FeatureRow;
2929
import java.util.Set;

ingestion/src/test/java/feast/NormalizeFeatureRows.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,7 @@
1919

2020
import com.google.common.collect.Lists;
2121
import com.google.common.primitives.UnsignedBytes;
22-
import feast.types.FeatureProto.Field;
22+
import feast.types.FieldProto.Field;
2323
import feast.types.FeatureRowProto.FeatureRow;
2424
import java.util.List;
2525
import org.apache.beam.sdk.transforms.MapElements;

ingestion/src/test/java/feast/ToOrderedFeatureRows.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,7 @@
1919

2020
import com.google.common.collect.Lists;
2121
import com.google.common.primitives.UnsignedBytes;
22-
import feast.types.FeatureProto.Field;
22+
import feast.types.FieldProto.Field;
2323
import feast.types.FeatureRowExtendedProto.FeatureRowExtended;
2424
import feast.types.FeatureRowProto.FeatureRow;
2525
import java.util.List;

protos/feast/types/Field.proto

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@ import "feast/types/Value.proto";
2121
package feast.types;
2222

2323
option java_package = "feast.types";
24-
option java_outer_classname = "FeatureProto";
24+
option java_outer_classname = "FieldProto";
2525
option go_package = "github.com/gojek/feast/protos/generated/go/feast/types";
2626

2727
message Field {

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

Lines changed: 44 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -37,14 +37,16 @@
3737
import feast.serving.ServingAPIProto.GetOnlineFeaturesResponse.FeatureDataSet;
3838
import feast.serving.exception.FeatureRetrievalException;
3939
import feast.storage.RedisProto.RedisKey;
40-
import feast.types.FeatureProto.Field;
40+
import feast.types.FieldProto.Field;
4141
import feast.types.FeatureRowProto.FeatureRow;
4242
import feast.types.FeatureRowProto.FeatureRow.Builder;
4343
import feast.types.ValueProto.Value;
4444
import io.opentracing.Scope;
4545
import io.opentracing.Tracer;
4646
import java.util.ArrayList;
47+
import java.util.HashMap;
4748
import java.util.List;
49+
import java.util.Map;
4850
import java.util.stream.Collectors;
4951
import lombok.extern.slf4j.Slf4j;
5052
import redis.clients.jedis.Jedis;
@@ -113,11 +115,10 @@ public GetOnlineFeaturesResponse getOnlineFeatures(GetFeaturesRequest request) {
113115
private List<RedisKey> getRedisKeys(List<String> entityNames,
114116
List<EntityDataSetRow> entityDataSetRows, FeatureSet featureSet) {
115117
try (Scope scope = tracer.buildSpan("Redis-makeRedisKeys").startActive(true)) {
116-
List<RedisKey> redisKeys = new ArrayList<>();
117118
String featureSetId = String.format("%s:%s", featureSet.getName(), featureSet.getVersion());
118-
for (EntityDataSetRow entityDataSetRow : entityDataSetRows) {
119-
redisKeys.add(makeRedisKey(featureSetId, entityNames, entityDataSetRow));
120-
}
119+
List<RedisKey> redisKeys = entityDataSetRows.parallelStream()
120+
.map(row -> makeRedisKey(featureSetId, entityNames, row))
121+
.collect(Collectors.toList());
121122
return redisKeys;
122123
}
123124
}
@@ -167,36 +168,51 @@ private List<FeatureRow> sendAndProcessMultiGet(List<RedisKey> redisKeys,
167168
List<String> requestedColumns, FeatureSet featureSet) throws InvalidProtocolBufferException {
168169
List<byte[]> jedisResps = sendMultiGet(redisKeys);
169170

170-
List<FeatureRow> featureRows = new ArrayList<>();
171-
172171
try (Scope scope = tracer.buildSpan("Redis-processResponse").startActive(true)) {
172+
List<FeatureRow> featureRows = new ArrayList<>();
173+
String featureSetName = String.format("%s:%s", featureSet.getName(), featureSet.getVersion());
173174
for (int i = 0; i < jedisResps.size(); i++) {
174-
byte[] jedisResp = jedisResps.get(i);
175-
if (jedisResp == null) {
176-
Builder emptyFeatureRowBuilder = FeatureRow.newBuilder()
177-
.setFeatureSet(String.format("%s:%s", featureSet.getName(), featureSet.getVersion()))
178-
.addAllFields(redisKeys.get(i).getEntitiesList())
179-
.setEventTimestamp(Timestamp.newBuilder().setSeconds(0).build());
180-
for (String requestedColumn : requestedColumns) {
181-
emptyFeatureRowBuilder.addFields(Field.newBuilder().setName(requestedColumn));
182-
}
183-
featureRows.add(emptyFeatureRowBuilder.build());
184-
} else {
185-
FeatureRow featureRow = FeatureRow.parseFrom(jedisResp);
186-
List<Field> fields = featureRow.getFieldsList().stream()
187-
.filter(f -> requestedColumns.contains(f.getName())).collect(Collectors.toList());
188-
featureRows.add(FeatureRow.newBuilder()
189-
.addAllFields(redisKeys.get(i).getEntitiesList())
190-
.addAllFields(fields)
191-
.setEventTimestamp(featureRow.getEventTimestamp())
192-
.setFeatureSet(String.format("%s:%s", featureSet.getName(), featureSet.getVersion()))
193-
.build());
194-
}
175+
featureRows.add(
176+
buildFeatureRow(jedisResps.get(i), featureSetName, redisKeys.get(i).getEntitiesList(),
177+
requestedColumns));
195178
}
196179
return featureRows;
197180
}
198181
}
199182

183+
/**
184+
* Build a featureRow given the request and the
185+
* @param jedisResponse
186+
* @param featureSet
187+
* @param entities
188+
* @param requestedColumns
189+
* @return
190+
* @throws InvalidProtocolBufferException
191+
*/
192+
private FeatureRow buildFeatureRow(byte[] jedisResponse, String featureSet, List<Field> entities,
193+
List<String> requestedColumns) throws InvalidProtocolBufferException {
194+
Builder featureRowBuilder = FeatureRow.newBuilder()
195+
.setFeatureSet(featureSet)
196+
.addAllFields(entities)
197+
.setEventTimestamp(Timestamp.newBuilder().setSeconds(0).build());
198+
199+
if (jedisResponse == null) {
200+
for (String requestedColumn : requestedColumns) {
201+
featureRowBuilder.addFields(Field.newBuilder().setName(requestedColumn));
202+
}
203+
return featureRowBuilder.build();
204+
}
205+
FeatureRow featureRow = FeatureRow.parseFrom(jedisResponse);
206+
List<Field> fields = featureRow.getFieldsList().stream()
207+
.filter(f -> requestedColumns.contains(f.getName())).collect(Collectors.toList());
208+
return featureRowBuilder
209+
.addAllFields(entities)
210+
.addAllFields(fields)
211+
.setEventTimestamp(featureRow.getEventTimestamp())
212+
.setFeatureSet(featureSet)
213+
.build();
214+
}
215+
200216
/**
201217
* Send a list of get request as an mget
202218
*

serving/src/test/java/feast/serving/service/RedisGrpcServingServiceTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -30,8 +30,8 @@
3030
import feast.serving.service.spec.SpecService;
3131
import feast.serving.testutil.FakeRedisCoreService;
3232
import feast.serving.testutil.RedisPopulator;
33-
import feast.types.FeatureProto.Field;
3433
import feast.types.FeatureRowProto.FeatureRow;
34+
import feast.types.FieldProto.Field;
3535
import feast.types.ValueProto.Value;
3636
import io.opentracing.util.GlobalTracer;
3737
import java.util.ArrayList;

serving/src/test/java/feast/serving/testutil/FeatureStoragePopulator.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,7 @@
1919

2020
import com.google.protobuf.ByteString;
2121
import feast.core.FeatureSetProto.FeatureSetSpec;
22-
import feast.types.FeatureProto.Field;
22+
import feast.types.FieldProto.Field;
2323
import feast.types.FeatureRowProto.FeatureRow;
2424
import feast.types.ValueProto.BoolList;
2525
import feast.types.ValueProto.BytesList;

serving/src/test/java/feast/serving/testutil/RedisPopulator.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@
2020
import feast.core.FeatureSetProto.EntitySpec;
2121
import feast.core.FeatureSetProto.FeatureSetSpec;
2222
import feast.storage.RedisProto.RedisKey;
23-
import feast.types.FeatureProto.Field;
23+
import feast.types.FieldProto.Field;
2424
import feast.types.FeatureRowProto.FeatureRow;
2525
import java.util.ArrayList;
2626
import java.util.List;

0 commit comments

Comments
 (0)