Skip to content

Commit bfac935

Browse files
author
zhilingc
committed
Implement API changes for BQ store
1 parent 5a9bda5 commit bfac935

9 files changed

Lines changed: 227 additions & 285 deletions

File tree

core/src/main/java/feast/core/grpc/CoreServiceImpl.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -135,6 +135,7 @@ public void applyFeatureSet(
135135
}
136136

137137
@Override
138+
@Transactional
138139
public void updateStore(UpdateStoreRequest request,
139140
StreamObserver<UpdateStoreResponse> responseObserver) {
140141
try {

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
package feast.serving.controller;
22

33
import feast.core.StoreProto.Store;
4-
import feast.serving.ServingAPIProto.GetFeastServingTypeRequest;
4+
import feast.serving.ServingAPIProto.GetFeastServingInfoRequest;
55
import feast.serving.service.CachedSpecService;
66
import feast.serving.service.ServingService;
77
import io.grpc.health.v1.HealthGrpc.HealthImplBase;
@@ -34,7 +34,7 @@ public void check(
3434

3535
try {
3636
Store store = specService.getStore();
37-
servingService.getFeastServingType(GetFeastServingTypeRequest.getDefaultInstance());
37+
servingService.getFeastServingInfo(GetFeastServingInfoRequest.getDefaultInstance());
3838
responseObserver.onNext(
3939
HealthCheckResponse.newBuilder().setStatus(ServingStatus.SERVING).build());
4040
} catch (Exception e) {

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

Lines changed: 15 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,13 @@
11
package feast.serving.controller;
22

33
import feast.serving.FeastProperties;
4+
import feast.serving.ServingAPIProto.GetBatchFeaturesRequest;
45
import feast.serving.ServingAPIProto.GetBatchFeaturesResponse;
5-
import feast.serving.ServingAPIProto.GetFeastServingTypeRequest;
6-
import feast.serving.ServingAPIProto.GetFeastServingTypeResponse;
7-
import feast.serving.ServingAPIProto.GetFeastServingVersionRequest;
8-
import feast.serving.ServingAPIProto.GetFeastServingVersionResponse;
9-
import feast.serving.ServingAPIProto.GetFeaturesRequest;
6+
import feast.serving.ServingAPIProto.GetFeastServingInfoRequest;
7+
import feast.serving.ServingAPIProto.GetFeastServingInfoResponse;
108
import feast.serving.ServingAPIProto.GetJobRequest;
119
import feast.serving.ServingAPIProto.GetJobResponse;
10+
import feast.serving.ServingAPIProto.GetOnlineFeaturesRequest;
1211
import feast.serving.ServingAPIProto.GetOnlineFeaturesResponse;
1312
import feast.serving.ServingServiceGrpc.ServingServiceImplBase;
1413
import feast.serving.service.ServingService;
@@ -37,28 +36,22 @@ public ServingServiceGRpcController(
3736
}
3837

3938
@Override
40-
public void getFeastServingVersion(
41-
GetFeastServingVersionRequest request,
42-
StreamObserver<GetFeastServingVersionResponse> responseObserver) {
43-
responseObserver.onNext(
44-
GetFeastServingVersionResponse.newBuilder().setVersion(version).build());
45-
responseObserver.onCompleted();
46-
}
47-
48-
@Override
49-
public void getFeastServingType(
50-
GetFeastServingTypeRequest request,
51-
StreamObserver<GetFeastServingTypeResponse> responseObserver) {
52-
responseObserver.onNext(servingService.getFeastServingType(request));
39+
public void getFeastServingInfo(GetFeastServingInfoRequest request,
40+
StreamObserver<GetFeastServingInfoResponse> responseObserver) {
41+
GetFeastServingInfoResponse feastServingInfo = servingService.getFeastServingInfo(request);
42+
feastServingInfo = feastServingInfo.toBuilder()
43+
.setVersion(version)
44+
.build();
45+
responseObserver.onNext(feastServingInfo);
5346
responseObserver.onCompleted();
5447
}
5548

5649
@Override
5750
public void getOnlineFeatures(
58-
GetFeaturesRequest request, StreamObserver<GetOnlineFeaturesResponse> responseObserver) {
51+
GetOnlineFeaturesRequest request, StreamObserver<GetOnlineFeaturesResponse> responseObserver) {
5952
Span span = tracer.buildSpan("getOnlineFeatures").start();
6053
try (Scope scope = tracer.scopeManager().activate(span, false)) {
61-
RequestHelper.validateRequest(request);
54+
RequestHelper.validateOnlineRequest(request);
6255
GetOnlineFeaturesResponse onlineFeatures = servingService.getOnlineFeatures(request);
6356
responseObserver.onNext(onlineFeatures);
6457
responseObserver.onCompleted();
@@ -70,8 +63,9 @@ public void getOnlineFeatures(
7063

7164
@Override
7265
public void getBatchFeatures(
73-
GetFeaturesRequest request, StreamObserver<GetBatchFeaturesResponse> responseObserver) {
66+
GetBatchFeaturesRequest request, StreamObserver<GetBatchFeaturesResponse> responseObserver) {
7467
try {
68+
RequestHelper.validateBatchRequest(request);
7569
GetBatchFeaturesResponse batchFeatures = servingService.getBatchFeatures(request);
7670
responseObserver.onNext(batchFeatures);
7771
responseObserver.onCompleted();

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

Lines changed: 11 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -3,8 +3,9 @@
33
import static feast.serving.util.mappers.ResponseJSONMapper.mapGetOnlineFeaturesResponse;
44

55
import feast.serving.FeastProperties;
6-
import feast.serving.ServingAPIProto.GetFeastServingVersionResponse;
7-
import feast.serving.ServingAPIProto.GetFeaturesRequest;
6+
import feast.serving.ServingAPIProto.GetFeastServingInfoRequest;
7+
import feast.serving.ServingAPIProto.GetFeastServingInfoResponse;
8+
import feast.serving.ServingAPIProto.GetOnlineFeaturesRequest;
89
import feast.serving.ServingAPIProto.GetOnlineFeaturesResponse;
910
import feast.serving.service.ServingService;
1011
import feast.serving.util.RequestHelper;
@@ -31,17 +32,20 @@ public ServingServiceRestController(
3132
this.tracer = tracer;
3233
}
3334

34-
@RequestMapping(value = "/api/v1/version", produces = "application/json")
35-
public GetFeastServingVersionResponse getVersion() {
36-
return GetFeastServingVersionResponse.newBuilder().setVersion(version).build();
35+
@RequestMapping(value = "/api/v1/info", produces = "application/json")
36+
public GetFeastServingInfoResponse getInfo() {
37+
GetFeastServingInfoResponse feastServingInfo = servingService
38+
.getFeastServingInfo(GetFeastServingInfoRequest.getDefaultInstance());
39+
return feastServingInfo.toBuilder().setVersion(version).build();
3740
}
3841

3942
@RequestMapping(
4043
value = "/api/v1/features/online",
4144
produces = "application/json",
4245
consumes = "application/json")
43-
public List<Map<String, Object>> getOnlineFeatures(@RequestBody GetFeaturesRequest request) {
44-
RequestHelper.validateRequest(request);
46+
public List<Map<String, Object>> getOnlineFeatures(
47+
@RequestBody GetOnlineFeaturesRequest request) {
48+
RequestHelper.validateOnlineRequest(request);
4549
GetOnlineFeaturesResponse onlineFeatures = servingService.getOnlineFeatures(request);
4650
return mapGetOnlineFeaturesResponse(onlineFeatures);
4751
}

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

Lines changed: 75 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -4,24 +4,33 @@
44
import com.google.cloud.bigquery.BigQueryException;
55
import com.google.cloud.bigquery.DatasetId;
66
import com.google.cloud.bigquery.ExtractJobConfiguration;
7+
import com.google.cloud.bigquery.Field;
8+
import com.google.cloud.bigquery.FormatOptions;
79
import com.google.cloud.bigquery.Job;
810
import com.google.cloud.bigquery.JobInfo;
11+
import com.google.cloud.bigquery.LoadJobConfiguration;
912
import com.google.cloud.bigquery.QueryJobConfiguration;
13+
import com.google.cloud.bigquery.Schema;
14+
import com.google.cloud.bigquery.Table;
15+
import com.google.cloud.bigquery.TableDefinition;
16+
import com.google.cloud.bigquery.TableId;
17+
import com.google.cloud.bigquery.TableInfo;
1018
import com.google.cloud.storage.Blob;
1119
import com.google.cloud.storage.Storage;
1220
import com.google.cloud.storage.Storage.BlobListOption;
1321
import com.google.common.collect.Lists;
1422
import feast.core.FeatureSetProto.FeatureSetSpec;
1523
import feast.serving.ServingAPIProto;
1624
import feast.serving.ServingAPIProto.DataFormat;
25+
import feast.serving.ServingAPIProto.DatasetSource;
1726
import feast.serving.ServingAPIProto.FeastServingType;
27+
import feast.serving.ServingAPIProto.GetBatchFeaturesRequest;
1828
import feast.serving.ServingAPIProto.GetBatchFeaturesResponse;
19-
import feast.serving.ServingAPIProto.GetFeastServingTypeRequest;
20-
import feast.serving.ServingAPIProto.GetFeastServingTypeResponse;
21-
import feast.serving.ServingAPIProto.GetFeaturesRequest;
22-
import feast.serving.ServingAPIProto.GetFeaturesRequest.EntityRow;
29+
import feast.serving.ServingAPIProto.GetFeastServingInfoRequest;
30+
import feast.serving.ServingAPIProto.GetFeastServingInfoResponse;
2331
import feast.serving.ServingAPIProto.GetJobRequest;
2432
import feast.serving.ServingAPIProto.GetJobResponse;
33+
import feast.serving.ServingAPIProto.GetOnlineFeaturesRequest;
2534
import feast.serving.ServingAPIProto.GetOnlineFeaturesResponse;
2635
import feast.serving.ServingAPIProto.JobStatus;
2736
import feast.serving.ServingAPIProto.JobType;
@@ -62,22 +71,31 @@ public BigQueryServingService(
6271
this.storage = storage;
6372
}
6473

74+
/**
75+
* {@inheritDoc}
76+
*/
6577
@Override
66-
public GetFeastServingTypeResponse getFeastServingType(
67-
GetFeastServingTypeRequest getFeastServingTypeRequest) {
68-
return GetFeastServingTypeResponse.newBuilder()
78+
public GetFeastServingInfoResponse getFeastServingInfo(
79+
GetFeastServingInfoRequest getFeastServingInfoRequest) {
80+
return GetFeastServingInfoResponse.newBuilder()
6981
.setType(FeastServingType.FEAST_SERVING_TYPE_BATCH)
82+
.setJobStagingLocation(jobStagingLocation)
7083
.build();
7184
}
7285

86+
/**
87+
* {@inheritDoc}
88+
*/
7389
@Override
74-
public GetOnlineFeaturesResponse getOnlineFeatures(GetFeaturesRequest getFeaturesRequest) {
90+
public GetOnlineFeaturesResponse getOnlineFeatures(GetOnlineFeaturesRequest getFeaturesRequest) {
7591
throw Status.UNIMPLEMENTED.withDescription("Method not implemented").asRuntimeException();
7692
}
7793

94+
/**
95+
* {@inheritDoc}
96+
*/
7897
@Override
79-
public GetBatchFeaturesResponse getBatchFeatures(GetFeaturesRequest getFeaturesRequest) {
80-
// TODO: Consider default maxAge in featureSetSpec during retrieval
98+
public GetBatchFeaturesResponse getBatchFeatures(GetBatchFeaturesRequest getFeaturesRequest) {
8199

82100
List<FeatureSetSpec> featureSetSpecs =
83101
getFeaturesRequest.getFeatureSetsList().stream()
@@ -93,30 +111,17 @@ public GetBatchFeaturesResponse getBatchFeatures(GetFeaturesRequest getFeaturesR
93111
.asRuntimeException();
94112
}
95113

96-
if (getFeaturesRequest.getEntityRowsCount() < 1) {
97-
throw Status.INVALID_ARGUMENT
98-
.withDescription(
99-
"entity_dataset_rows is required for batch retrieval in order to filter the retrieved entities.")
100-
.asRuntimeException();
101-
}
102-
103-
for (EntityRow entityRow :
104-
getFeaturesRequest.getEntityRowsList()) {
105-
if (entityRow.getEntityTimestamp().getSeconds() == 0) {
106-
throw Status.INVALID_ARGUMENT
107-
.withDescription(
108-
"entity_timestamp field in entity_dataset_row is required for batch retrieval.")
109-
.asRuntimeException();
110-
}
111-
}
114+
Table entityTable = loadEntities(getFeaturesRequest.getDatasetSource());
115+
Schema entityTableSchema = entityTable.getDefinition().getSchema();
116+
List<String> entityNames = entityTableSchema.getFields().stream().map(Field::getName)
117+
.collect(Collectors.toList());
112118

113119
final String query =
114120
BigQueryUtil.createQuery(
115121
getFeaturesRequest.getFeatureSetsList(),
116122
featureSetSpecs,
117-
Lists.newArrayList(getFeaturesRequest.getEntityRows(0).getFieldsMap().keySet()),
118-
getFeaturesRequest.getEntityRowsList(),
119-
datasetId);
123+
entityNames,
124+
datasetId, entityTable.getFriendlyName());
120125
log.debug("Running BigQuery query: {}", query);
121126

122127
String feastJobId = UUID.randomUUID().toString();
@@ -210,6 +215,9 @@ public GetBatchFeaturesResponse getBatchFeatures(GetFeaturesRequest getFeaturesR
210215
return GetBatchFeaturesResponse.newBuilder().setJob(feastJob).build();
211216
}
212217

218+
/**
219+
* {@inheritDoc}
220+
*/
213221
@Override
214222
public GetJobResponse getJob(GetJobRequest getJobRequest) {
215223
Optional<ServingAPIProto.Job> job = jobService.get(getJobRequest.getJob().getId());
@@ -220,4 +228,42 @@ public GetJobResponse getJob(GetJobRequest getJobRequest) {
220228
}
221229
return GetJobResponse.newBuilder().setJob(job.get()).build();
222230
}
231+
232+
private Table loadEntities(DatasetSource datasetSource) {
233+
switch (datasetSource.getDatasetSourceCase()) {
234+
case FILE_SOURCE:
235+
String tableName = generateTemporaryTableName();
236+
TableId tableId = TableId.of(projectId, datasetId, tableName);
237+
// Currently only avro supported
238+
if (datasetSource.getFileSource().getDataFormat() != DataFormat.DATA_FORMAT_AVRO) {
239+
throw Status.INVALID_ARGUMENT
240+
.withDescription("Invalid file format, only avro supported")
241+
.asRuntimeException();
242+
}
243+
LoadJobConfiguration loadJobConfiguration = LoadJobConfiguration.of(tableId,
244+
datasetSource.getFileSource().getFileUrisList(),
245+
FormatOptions.avro());
246+
Job job = bigquery.create(JobInfo.of(loadJobConfiguration));
247+
try {
248+
job.waitFor();
249+
return bigquery.getTable(tableId);
250+
} catch (InterruptedException e) {
251+
throw Status.INTERNAL
252+
.withDescription("Failed to load entity dataset into store")
253+
.withCause(e)
254+
.asRuntimeException();
255+
}
256+
case DATASETSOURCE_NOT_SET:
257+
default:
258+
throw Status.INVALID_ARGUMENT
259+
.withDescription("Data source must be set.")
260+
.asRuntimeException();
261+
}
262+
}
263+
264+
private String generateTemporaryTableName() {
265+
String source = String.format("feast_serving_%d", System.currentTimeMillis());
266+
UUID uuid = UUID.fromString(source);
267+
return uuid.toString().replaceAll("-", "_");
268+
}
223269
}

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

Lines changed: 11 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -20,19 +20,18 @@
2020
import com.google.protobuf.AbstractMessageLite;
2121
import com.google.protobuf.Duration;
2222
import com.google.protobuf.InvalidProtocolBufferException;
23-
import feast.core.CoreServiceProto.GetFeatureSetsRequest;
24-
import feast.core.CoreServiceProto.GetFeatureSetsRequest.Filter;
2523
import feast.core.FeatureSetProto.EntitySpec;
2624
import feast.core.FeatureSetProto.FeatureSetSpec;
2725
import feast.serving.ServingAPIProto.FeastServingType;
26+
import feast.serving.ServingAPIProto.GetBatchFeaturesRequest;
2827
import feast.serving.ServingAPIProto.GetBatchFeaturesResponse;
29-
import feast.serving.ServingAPIProto.GetFeastServingTypeRequest;
30-
import feast.serving.ServingAPIProto.GetFeastServingTypeResponse;
31-
import feast.serving.ServingAPIProto.GetFeaturesRequest;
32-
import feast.serving.ServingAPIProto.GetFeaturesRequest.EntityRow;
33-
import feast.serving.ServingAPIProto.GetFeaturesRequest.FeatureSet;
28+
import feast.serving.ServingAPIProto.GetFeastServingInfoRequest;
29+
import feast.serving.ServingAPIProto.GetFeastServingInfoResponse;
3430
import feast.serving.ServingAPIProto.GetJobRequest;
3531
import feast.serving.ServingAPIProto.GetJobResponse;
32+
import feast.serving.ServingAPIProto.GetOnlineFeaturesRequest;
33+
import feast.serving.ServingAPIProto.GetOnlineFeaturesRequest.EntityRow;
34+
import feast.serving.ServingAPIProto.GetOnlineFeaturesRequest.FeatureSet;
3635
import feast.serving.ServingAPIProto.GetOnlineFeaturesResponse;
3736
import feast.serving.ServingAPIProto.GetOnlineFeaturesResponse.FieldValues;
3837
import feast.storage.RedisProto.RedisKey;
@@ -64,16 +63,16 @@ public RedisServingService(JedisPool jedisPool, CachedSpecService specService, T
6463

6564
/** {@inheritDoc} */
6665
@Override
67-
public GetFeastServingTypeResponse getFeastServingType(
68-
GetFeastServingTypeRequest getFeastServingTypeRequest) {
69-
return GetFeastServingTypeResponse.newBuilder()
66+
public GetFeastServingInfoResponse getFeastServingInfo(
67+
GetFeastServingInfoRequest getFeastServingInfoRequest) {
68+
return GetFeastServingInfoResponse.newBuilder()
7069
.setType(FeastServingType.FEAST_SERVING_TYPE_ONLINE)
7170
.build();
7271
}
7372

7473
/** {@inheritDoc} */
7574
@Override
76-
public GetOnlineFeaturesResponse getOnlineFeatures(GetFeaturesRequest request) {
75+
public GetOnlineFeaturesResponse getOnlineFeatures(GetOnlineFeaturesRequest request) {
7776
try (Scope scope = tracer.buildSpan("Redis-getOnlineFeatures").startActive(true)) {
7877
GetOnlineFeaturesResponse.Builder getOnlineFeaturesResponseBuilder =
7978
GetOnlineFeaturesResponse.newBuilder();
@@ -120,7 +119,7 @@ public GetOnlineFeaturesResponse getOnlineFeatures(GetFeaturesRequest request) {
120119
}
121120

122121
@Override
123-
public GetBatchFeaturesResponse getBatchFeatures(GetFeaturesRequest getFeaturesRequest) {
122+
public GetBatchFeaturesResponse getBatchFeatures(GetBatchFeaturesRequest getFeaturesRequest) {
124123
throw Status.UNIMPLEMENTED.withDescription("Method not implemented").asRuntimeException();
125124
}
126125

Lines changed: 8 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1,20 +1,21 @@
11
package feast.serving.service;
22

3+
import feast.serving.ServingAPIProto.GetBatchFeaturesRequest;
34
import feast.serving.ServingAPIProto.GetBatchFeaturesResponse;
4-
import feast.serving.ServingAPIProto.GetFeastServingTypeRequest;
5-
import feast.serving.ServingAPIProto.GetFeastServingTypeResponse;
6-
import feast.serving.ServingAPIProto.GetFeaturesRequest;
5+
import feast.serving.ServingAPIProto.GetFeastServingInfoRequest;
6+
import feast.serving.ServingAPIProto.GetFeastServingInfoResponse;
77
import feast.serving.ServingAPIProto.GetJobRequest;
88
import feast.serving.ServingAPIProto.GetJobResponse;
9+
import feast.serving.ServingAPIProto.GetOnlineFeaturesRequest;
910
import feast.serving.ServingAPIProto.GetOnlineFeaturesResponse;
1011

1112
public interface ServingService {
12-
GetFeastServingTypeResponse getFeastServingType(
13-
GetFeastServingTypeRequest getFeastServingTypeRequest);
13+
GetFeastServingInfoResponse getFeastServingInfo(
14+
GetFeastServingInfoRequest getFeastServingInfoRequest);
1415

15-
GetOnlineFeaturesResponse getOnlineFeatures(GetFeaturesRequest getFeaturesRequest);
16+
GetOnlineFeaturesResponse getOnlineFeatures(GetOnlineFeaturesRequest getFeaturesRequest);
1617

17-
GetBatchFeaturesResponse getBatchFeatures(GetFeaturesRequest getFeaturesRequest);
18+
GetBatchFeaturesResponse getBatchFeatures(GetBatchFeaturesRequest getFeaturesRequest);
1819

1920
GetJobResponse getJob(GetJobRequest getJobRequest);
2021
}

0 commit comments

Comments
 (0)