From 375071982b71cca7b500c386ae25830cc0313c0a Mon Sep 17 00:00:00 2001 From: zhilingc Date: Wed, 25 Mar 2020 16:09:35 +0800 Subject: [PATCH 1/5] API and docstring tweaks --- .../main/java/feast/ingestion/ImportJob.java | 2 +- .../configuration/ServingServiceConfig.java | 16 ++- ...ice.java => HistoricalServingService.java} | 63 ++++++++--- .../retrieval/HistoricalRetrievalResult.java | 100 ++++++++++++++++++ ...etriever.java => HistoricalRetriever.java} | 22 ++-- .../api/retrieval/OnlineRetriever.java | 5 +- .../retrieval/OnlineRetrieverResponse.java | 44 -------- .../feast/storage/api/write/FeatureSink.java | 2 +- ....java => BigQueryHistoricalRetriever.java} | 53 ++++------ .../bigquery/retrieval/SubqueryCallable.java | 2 +- .../bigquery/write/BigQueryFeatureSink.java | 2 +- storage/connectors/pom.xml | 9 +- .../redis/write/RedisFeatureSink.java | 2 +- .../redis/write/RedisFeatureSinkTest.java | 14 +-- 14 files changed, 207 insertions(+), 129 deletions(-) rename serving/src/main/java/feast/serving/service/{BatchServingService.java => HistoricalServingService.java} (55%) create mode 100644 storage/api/src/main/java/feast/storage/api/retrieval/HistoricalRetrievalResult.java rename storage/api/src/main/java/feast/storage/api/retrieval/{BatchRetriever.java => HistoricalRetriever.java} (57%) delete mode 100644 storage/api/src/main/java/feast/storage/api/retrieval/OnlineRetrieverResponse.java rename storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/retrieval/{BigQueryBatchRetriever.java => BigQueryHistoricalRetriever.java} (89%) diff --git a/ingestion/src/main/java/feast/ingestion/ImportJob.java b/ingestion/src/main/java/feast/ingestion/ImportJob.java index c8982ca24ad..5b0bf88924a 100644 --- a/ingestion/src/main/java/feast/ingestion/ImportJob.java +++ b/ingestion/src/main/java/feast/ingestion/ImportJob.java @@ -137,7 +137,7 @@ public static PipelineResult runPipeline(ImportOptions options) throws IOExcepti // Step 3. Write FeatureRow to the corresponding Store. WriteResult writeFeatureRows = - validatedRows.get(FEATURE_ROW_OUT).apply("WriteFeatureRowToStore", featureSink.write()); + validatedRows.get(FEATURE_ROW_OUT).apply("WriteFeatureRowToStore", featureSink.writer()); // Step 4. Write FailedElements to a dead letter table in BigQuery. if (options.getDeadLetterTableSpec() != null) { diff --git a/serving/src/main/java/feast/serving/configuration/ServingServiceConfig.java b/serving/src/main/java/feast/serving/configuration/ServingServiceConfig.java index e9819c275fe..a149b3fa14c 100644 --- a/serving/src/main/java/feast/serving/configuration/ServingServiceConfig.java +++ b/serving/src/main/java/feast/serving/configuration/ServingServiceConfig.java @@ -25,15 +25,11 @@ import feast.core.StoreProto.Store.RedisConfig; import feast.core.StoreProto.Store.Subscription; import feast.serving.FeastProperties; -import feast.serving.service.BatchServingService; -import feast.serving.service.JobService; -import feast.serving.service.NoopJobService; -import feast.serving.service.OnlineServingService; -import feast.serving.service.ServingService; +import feast.serving.service.*; import feast.serving.specs.CachedSpecService; -import feast.storage.api.retrieval.BatchRetriever; +import feast.storage.api.retrieval.HistoricalRetriever; import feast.storage.api.retrieval.OnlineRetriever; -import feast.storage.connectors.bigquery.retrieval.BigQueryBatchRetriever; +import feast.storage.connectors.bigquery.retrieval.BigQueryHistoricalRetriever; import feast.storage.connectors.redis.retrieval.RedisOnlineRetriever; import io.opentracing.Tracer; import java.util.Map; @@ -109,8 +105,8 @@ public ServingService servingService( "Unable to instantiate jobService for BigQuery store."); } - BatchRetriever bqRetriever = - BigQueryBatchRetriever.builder() + HistoricalRetriever bqRetriever = + BigQueryHistoricalRetriever.builder() .setBigquery(bigquery) .setDatasetId(bqConfig.getDatasetId()) .setProjectId(bqConfig.getProjectId()) @@ -121,7 +117,7 @@ public ServingService servingService( .setStorage(storage) .build(); - servingService = new BatchServingService(bqRetriever, specService, jobService); + servingService = new HistoricalServingService(bqRetriever, specService, jobService); break; case CASSANDRA: case UNRECOGNIZED: diff --git a/serving/src/main/java/feast/serving/service/BatchServingService.java b/serving/src/main/java/feast/serving/service/HistoricalServingService.java similarity index 55% rename from serving/src/main/java/feast/serving/service/BatchServingService.java rename to serving/src/main/java/feast/serving/service/HistoricalServingService.java index b33a9b4ad7a..7d389a9e12d 100644 --- a/serving/src/main/java/feast/serving/service/BatchServingService.java +++ b/serving/src/main/java/feast/serving/service/HistoricalServingService.java @@ -18,25 +18,29 @@ import feast.serving.ServingAPIProto; import feast.serving.ServingAPIProto.*; +import feast.serving.ServingAPIProto.Job.Builder; import feast.serving.specs.CachedSpecService; -import feast.storage.api.retrieval.BatchRetriever; import feast.storage.api.retrieval.FeatureSetRequest; -import feast.storage.connectors.bigquery.retrieval.BigQueryBatchRetriever; +import feast.storage.api.retrieval.HistoricalRetrievalResult; +import feast.storage.api.retrieval.HistoricalRetriever; +import feast.storage.connectors.bigquery.retrieval.BigQueryHistoricalRetriever; import io.grpc.Status; import java.util.List; import java.util.Optional; +import java.util.UUID; import org.slf4j.Logger; -public class BatchServingService implements ServingService { +public class HistoricalServingService implements ServingService { - private static final Logger log = org.slf4j.LoggerFactory.getLogger(BatchServingService.class); + private static final Logger log = + org.slf4j.LoggerFactory.getLogger(HistoricalServingService.class); - private final BatchRetriever retriever; + private final HistoricalRetriever retriever; private final CachedSpecService specService; private final JobService jobService; - public BatchServingService( - BatchRetriever retriever, CachedSpecService specService, JobService jobService) { + public HistoricalServingService( + HistoricalRetriever retriever, CachedSpecService specService, JobService jobService) { this.retriever = retriever; this.specService = specService; this.jobService = jobService; @@ -47,10 +51,11 @@ public BatchServingService( public GetFeastServingInfoResponse getFeastServingInfo( GetFeastServingInfoRequest getFeastServingInfoRequest) { try { - BigQueryBatchRetriever bigQueryBatchRetriever = (BigQueryBatchRetriever) retriever; + BigQueryHistoricalRetriever bigQueryHistoricalRetriever = + (BigQueryHistoricalRetriever) retriever; return GetFeastServingInfoResponse.newBuilder() .setType(FeastServingType.FEAST_SERVING_TYPE_BATCH) - .setJobStagingLocation(bigQueryBatchRetriever.jobStagingLocation()) + .setJobStagingLocation(bigQueryHistoricalRetriever.jobStagingLocation()) .build(); } catch (Exception e) { return GetFeastServingInfoResponse.newBuilder() @@ -70,9 +75,28 @@ public GetOnlineFeaturesResponse getOnlineFeatures(GetOnlineFeaturesRequest getF public GetBatchFeaturesResponse getBatchFeatures(GetBatchFeaturesRequest getFeaturesRequest) { List featureSetRequests = specService.getFeatureSets(getFeaturesRequest.getFeaturesList()); - Job feastJob = retriever.getBatchFeatures(getFeaturesRequest, featureSetRequests); - jobService.upsert(feastJob); - return GetBatchFeaturesResponse.newBuilder().setJob(feastJob).build(); + String retrievalId = UUID.randomUUID().toString(); + Job runningJob = + Job.newBuilder() + .setId(retrievalId) + .setType(JobType.JOB_TYPE_DOWNLOAD) + .setStatus(JobStatus.JOB_STATUS_RUNNING) + .build(); + jobService.upsert(runningJob); + Thread thread = + new Thread( + new Runnable() { + @Override + public void run() { + HistoricalRetrievalResult result = + retriever.getHistoricalFeatures( + retrievalId, getFeaturesRequest.getDatasetSource(), featureSetRequests); + jobService.upsert(resultToJob(result)); + } + }); + thread.start(); + + return GetBatchFeaturesResponse.newBuilder().setJob(runningJob).build(); } /** {@inheritDoc} */ @@ -86,4 +110,19 @@ public GetJobResponse getJob(GetJobRequest getJobRequest) { } return GetJobResponse.newBuilder().setJob(job.get()).build(); } + + private Job resultToJob(HistoricalRetrievalResult result) { + Builder builder = + Job.newBuilder() + .setId(result.getId()) + .setType(JobType.JOB_TYPE_DOWNLOAD) + .setStatus(result.getStatus()); + if (result.hasError()) { + return builder.setError(result.getError()).build(); + } + return builder + .addAllFileUris(result.getFileUris()) + .setDataFormat(result.getDataFormat()) + .build(); + } } diff --git a/storage/api/src/main/java/feast/storage/api/retrieval/HistoricalRetrievalResult.java b/storage/api/src/main/java/feast/storage/api/retrieval/HistoricalRetrievalResult.java new file mode 100644 index 00000000000..59e08dbe8c7 --- /dev/null +++ b/storage/api/src/main/java/feast/storage/api/retrieval/HistoricalRetrievalResult.java @@ -0,0 +1,100 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * Copyright 2018-2020 The Feast Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package feast.storage.api.retrieval; + +import com.google.auto.value.AutoValue; +import feast.serving.ServingAPIProto.DataFormat; +import feast.serving.ServingAPIProto.JobStatus; +import java.io.Serializable; +import java.util.List; +import javax.annotation.Nullable; + +/** Result of a historical feature retrieval request. */ +@AutoValue +public abstract class HistoricalRetrievalResult implements Serializable { + + public abstract String getId(); + + public abstract JobStatus getStatus(); + + @Nullable + public abstract String getError(); + + @Nullable + public abstract List getFileUris(); + + @Nullable + public abstract DataFormat getDataFormat(); + + /** + * Instantiates a {@link HistoricalRetrievalResult} indicating that the retrieval was a failure, + * together with its associated error. + * + * @param id retrieval id identifying the retrieval request. + * @param error error that occurred + * @return {@link HistoricalRetrievalResult} + */ + public static HistoricalRetrievalResult errorResult(String id, Exception error) { + return newBuilder() + .setId(id) + .setStatus(JobStatus.JOB_STATUS_DONE) + .setError(error.getMessage()) + .build(); + } + + /** + * Instantiates a {@link HistoricalRetrievalResult} indicating that the retrieval was a success, + * together with the location of the output. + * + * @param id retrieval id identifying the retrieval request + * @param fileUris list of output file URIs + * @param dataFormat data format of the output files + * @return + */ + public static HistoricalRetrievalResult successResult( + String id, List fileUris, DataFormat dataFormat) { + return newBuilder() + .setId(id) + .setStatus(JobStatus.JOB_STATUS_DONE) + .setFileUris(fileUris) + .setDataFormat(dataFormat) + .build(); + } + + static Builder newBuilder() { + return new AutoValue_HistoricalRetrievalResult.Builder(); + } + + @AutoValue.Builder + abstract static class Builder { + abstract Builder setId(String id); + + abstract Builder setStatus(JobStatus jobStatus); + + abstract Builder setError(String error); + + abstract Builder setFileUris(List fileUris); + + abstract Builder setDataFormat(DataFormat dataFormat); + + abstract HistoricalRetrievalResult build(); + } + + public boolean hasError() { + return getError() != null; + } +} diff --git a/storage/api/src/main/java/feast/storage/api/retrieval/BatchRetriever.java b/storage/api/src/main/java/feast/storage/api/retrieval/HistoricalRetriever.java similarity index 57% rename from storage/api/src/main/java/feast/storage/api/retrieval/BatchRetriever.java rename to storage/api/src/main/java/feast/storage/api/retrieval/HistoricalRetriever.java index 79bd6ad3281..3533ed140f8 100644 --- a/storage/api/src/main/java/feast/storage/api/retrieval/BatchRetriever.java +++ b/storage/api/src/main/java/feast/storage/api/retrieval/HistoricalRetriever.java @@ -16,22 +16,26 @@ */ package feast.storage.api.retrieval; -import feast.serving.ServingAPIProto; +import feast.serving.ServingAPIProto.DatasetSource; import java.util.List; -/** Interface for implementing user defined retrieval functionality from Batch/historical stores. */ -public interface BatchRetriever { +/** + * A historical retriever is a feature retriever that retrieves feature data corresponding to + * provided entities over a given period of time. + */ +public interface HistoricalRetriever { /** * Get all features corresponding to the provided batch features request. * - * @param request {@link ServingAPIProto.GetBatchFeaturesRequest} containing requested features - * and file containing entity columns. + * @param retrievalId String that uniquely identifies this retrieval request. + * @param datasetSource {@link DatasetSource} containing source to load the dataset containing + * entity columns. * @param featureSetRequests List of {@link FeatureSetRequest} to feature references in the * request tied to that feature set. - * @return {@link ServingAPIProto.Job} if successful, contains the location of the results, else - * contains the error to be returned to the user. + * @return {@link HistoricalRetrievalResult} if successful, contains the location of the results, + * else contains the error to be returned to the user. */ - ServingAPIProto.Job getBatchFeatures( - ServingAPIProto.GetBatchFeaturesRequest request, List featureSetRequests); + HistoricalRetrievalResult getHistoricalFeatures( + String retrievalId, DatasetSource datasetSource, List featureSetRequests); } diff --git a/storage/api/src/main/java/feast/storage/api/retrieval/OnlineRetriever.java b/storage/api/src/main/java/feast/storage/api/retrieval/OnlineRetriever.java index c4c80f0b093..72888bfc081 100644 --- a/storage/api/src/main/java/feast/storage/api/retrieval/OnlineRetriever.java +++ b/storage/api/src/main/java/feast/storage/api/retrieval/OnlineRetriever.java @@ -20,7 +20,10 @@ import feast.types.FeatureRowProto.FeatureRow; import java.util.List; -/** Interface for implementing user defined retrieval functionality from Online stores. */ +/** + * An online retriever is a feature retriever that retrieves the latest feature data corresponding + * to provided entities. + */ public interface OnlineRetriever { /** diff --git a/storage/api/src/main/java/feast/storage/api/retrieval/OnlineRetrieverResponse.java b/storage/api/src/main/java/feast/storage/api/retrieval/OnlineRetrieverResponse.java deleted file mode 100644 index 5d789441266..00000000000 --- a/storage/api/src/main/java/feast/storage/api/retrieval/OnlineRetrieverResponse.java +++ /dev/null @@ -1,44 +0,0 @@ -/* - * SPDX-License-Identifier: Apache-2.0 - * Copyright 2018-2020 The Feast Authors - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package feast.storage.api.retrieval; - -import feast.types.FeatureRowProto; - -/** Response from an online store. */ -public interface OnlineRetrieverResponse { - - /** - * Checks whether the response is empty, i.e. feature does not exist in the store - * - * @return boolean - */ - boolean isEmpty(); - - /** - * Get the featureset associated with this response. - * - * @return String featureset reference in format featureSet:version - */ - String getFeatureSet(); - - /** - * Parse response to FeatureRow - * - * @return {@link FeatureRowProto.FeatureRow} - */ - FeatureRowProto.FeatureRow toFeatureRow(); -} diff --git a/storage/api/src/main/java/feast/storage/api/write/FeatureSink.java b/storage/api/src/main/java/feast/storage/api/write/FeatureSink.java index cc93e0f6399..9fa41d0c204 100644 --- a/storage/api/src/main/java/feast/storage/api/write/FeatureSink.java +++ b/storage/api/src/main/java/feast/storage/api/write/FeatureSink.java @@ -50,5 +50,5 @@ public interface FeatureSink extends Serializable { * * @return {@link PTransform} */ - PTransform, WriteResult> write(); + PTransform, WriteResult> writer(); } diff --git a/storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/retrieval/BigQueryBatchRetriever.java b/storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/retrieval/BigQueryHistoricalRetriever.java similarity index 89% rename from storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/retrieval/BigQueryBatchRetriever.java rename to storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/retrieval/BigQueryHistoricalRetriever.java index 881cbc18ebb..46d4fea11c7 100644 --- a/storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/retrieval/BigQueryBatchRetriever.java +++ b/storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/retrieval/BigQueryHistoricalRetriever.java @@ -25,8 +25,10 @@ import com.google.cloud.storage.Blob; import com.google.cloud.storage.Storage; import feast.serving.ServingAPIProto; -import feast.storage.api.retrieval.BatchRetriever; +import feast.serving.ServingAPIProto.DatasetSource; import feast.storage.api.retrieval.FeatureSetRequest; +import feast.storage.api.retrieval.HistoricalRetrievalResult; +import feast.storage.api.retrieval.HistoricalRetriever; import io.grpc.Status; import java.io.IOException; import java.util.ArrayList; @@ -38,9 +40,10 @@ import org.threeten.bp.Duration; @AutoValue -public abstract class BigQueryBatchRetriever implements BatchRetriever { +public abstract class BigQueryHistoricalRetriever implements HistoricalRetriever { - private static final Logger log = org.slf4j.LoggerFactory.getLogger(BigQueryBatchRetriever.class); + private static final Logger log = + org.slf4j.LoggerFactory.getLogger(BigQueryHistoricalRetriever.class); public static final long TEMP_TABLE_EXPIRY_DURATION_MS = Duration.ofDays(1).toMillis(); private static final long SUBQUERY_TIMEOUT_SECS = 900; // 15 minutes @@ -60,7 +63,7 @@ public abstract class BigQueryBatchRetriever implements BatchRetriever { public abstract Storage storage(); public static Builder builder() { - return new AutoValue_BigQueryBatchRetriever.Builder(); + return new AutoValue_BigQueryHistoricalRetriever.Builder(); } @AutoValue.Builder @@ -79,15 +82,12 @@ public abstract static class Builder { public abstract Builder setStorage(Storage storage); - public abstract BigQueryBatchRetriever build(); + public abstract BigQueryHistoricalRetriever build(); } @Override - public ServingAPIProto.Job getBatchFeatures( - ServingAPIProto.GetBatchFeaturesRequest request, List featureSetRequests) { - // 0. Generate job ID - String feastJobId = UUID.randomUUID().toString(); - + public HistoricalRetrievalResult getHistoricalFeatures( + String retrievalId, DatasetSource datasetSource, List featureSetRequests) { List featureSetQueryInfos = QueryTemplater.getFeatureSetInfos(featureSetRequests); @@ -95,18 +95,15 @@ public ServingAPIProto.Job getBatchFeatures( Table entityTable; String entityTableName; try { - entityTable = loadEntities(request.getDatasetSource()); + entityTable = loadEntities(datasetSource); TableId entityTableWithUUIDs = generateUUIDs(entityTable); entityTableName = generateFullTableName(entityTableWithUUIDs); } catch (Exception e) { - return ServingAPIProto.Job.newBuilder() - .setId(feastJobId) - .setType(ServingAPIProto.JobType.JOB_TYPE_DOWNLOAD) - .setStatus(ServingAPIProto.JobStatus.JOB_STATUS_DONE) - .setDataFormat(ServingAPIProto.DataFormat.DATA_FORMAT_AVRO) - .setError(String.format("Unable to load entity table to BigQuery: %s", e.toString())) - .build(); + return HistoricalRetrievalResult.errorResult( + retrievalId, + new RuntimeException( + String.format("Unable to load entity table to BigQuery: %s", e.toString()))); } Schema entityTableSchema = entityTable.getDefinition().getSchema(); @@ -132,7 +129,7 @@ public ServingAPIProto.Job getBatchFeatures( entityTableName, entityTableColumnNames, featureSetQueryInfos, featureSetQueries); queryConfig = queryJob.getConfiguration(); String exportTableDestinationUri = - String.format("%s/%s/*.avro", jobStagingLocation(), feastJobId); + String.format("%s/%s/*.avro", jobStagingLocation(), retrievalId); // 5. Export the table // Hardcode the format to Avro for now @@ -143,23 +140,13 @@ public ServingAPIProto.Job getBatchFeatures( waitForJob(extractJob); } catch (BigQueryException | InterruptedException | IOException e) { - return ServingAPIProto.Job.newBuilder() - .setId(feastJobId) - .setType(ServingAPIProto.JobType.JOB_TYPE_DOWNLOAD) - .setStatus(ServingAPIProto.JobStatus.JOB_STATUS_DONE) - .setError(e.getMessage()) - .build(); + return HistoricalRetrievalResult.errorResult(retrievalId, e); } - List fileUris = parseOutputFileURIs(feastJobId); + List fileUris = parseOutputFileURIs(retrievalId); - return ServingAPIProto.Job.newBuilder() - .setId(feastJobId) - .setType(ServingAPIProto.JobType.JOB_TYPE_DOWNLOAD) - .setStatus(ServingAPIProto.JobStatus.JOB_STATUS_DONE) - .addAllFileUris(fileUris) - .setDataFormat(ServingAPIProto.DataFormat.DATA_FORMAT_AVRO) - .build(); + return HistoricalRetrievalResult.successResult( + retrievalId, fileUris, ServingAPIProto.DataFormat.DATA_FORMAT_AVRO); } private TableId generateUUIDs(Table loadedEntityTable) { diff --git a/storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/retrieval/SubqueryCallable.java b/storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/retrieval/SubqueryCallable.java index 4a36da21223..651fe8bf977 100644 --- a/storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/retrieval/SubqueryCallable.java +++ b/storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/retrieval/SubqueryCallable.java @@ -16,7 +16,7 @@ */ package feast.storage.connectors.bigquery.retrieval; -import static feast.storage.connectors.bigquery.retrieval.BigQueryBatchRetriever.TEMP_TABLE_EXPIRY_DURATION_MS; +import static feast.storage.connectors.bigquery.retrieval.BigQueryHistoricalRetriever.TEMP_TABLE_EXPIRY_DURATION_MS; import static feast.storage.connectors.bigquery.retrieval.QueryTemplater.generateFullTableName; import com.google.auto.value.AutoValue; diff --git a/storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/write/BigQueryFeatureSink.java b/storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/write/BigQueryFeatureSink.java index b38728fa966..828348859bf 100644 --- a/storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/write/BigQueryFeatureSink.java +++ b/storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/write/BigQueryFeatureSink.java @@ -125,7 +125,7 @@ public void prepareWrite(FeatureSetProto.FeatureSet featureSet) { } @Override - public PTransform, WriteResult> write() { + public PTransform, WriteResult> writer() { return new BigQueryWrite(DatasetId.of(getProjectId(), getDatasetId())); } diff --git a/storage/connectors/pom.xml b/storage/connectors/pom.xml index 2cb0b106081..b52668a31a4 100644 --- a/storage/connectors/pom.xml +++ b/storage/connectors/pom.xml @@ -1,6 +1,5 @@ - pom dev.feast @@ -11,6 +10,7 @@ 4.0.0 feast-storage-connectors + pom Feast Storage Connectors @@ -46,13 +46,6 @@ feast-storage-api ${project.version} - - - junit - junit - 4.12 - test - diff --git a/storage/connectors/redis/src/main/java/feast/storage/connectors/redis/write/RedisFeatureSink.java b/storage/connectors/redis/src/main/java/feast/storage/connectors/redis/write/RedisFeatureSink.java index 2f566bb78c2..51c8a8d58bc 100644 --- a/storage/connectors/redis/src/main/java/feast/storage/connectors/redis/write/RedisFeatureSink.java +++ b/storage/connectors/redis/src/main/java/feast/storage/connectors/redis/write/RedisFeatureSink.java @@ -68,7 +68,7 @@ public void prepareWrite(FeatureSet featureSet) { } @Override - public PTransform, WriteResult> write() { + public PTransform, WriteResult> writer() { return new RedisCustomIO.Write(getRedisConfig(), getFeatureSetSpecs()); } } diff --git a/storage/connectors/redis/src/test/java/feast/storage/connectors/redis/write/RedisFeatureSinkTest.java b/storage/connectors/redis/src/test/java/feast/storage/connectors/redis/write/RedisFeatureSinkTest.java index ddabed8fadf..9b0fb4014c6 100644 --- a/storage/connectors/redis/src/test/java/feast/storage/connectors/redis/write/RedisFeatureSinkTest.java +++ b/storage/connectors/redis/src/test/java/feast/storage/connectors/redis/write/RedisFeatureSinkTest.java @@ -157,7 +157,7 @@ public void shouldWriteToRedis() { .addFields(field("feature", "two", Enum.STRING)) .build()); - p.apply(Create.of(featureRows)).apply(redisFeatureSink.write()); + p.apply(Create.of(featureRows)).apply(redisFeatureSink.writer()); p.run(); kvs.forEach( @@ -199,7 +199,7 @@ public void shouldRetryFailConnection() throws InterruptedException { PCollection failedElementCount = p.apply(Create.of(featureRows)) - .apply(redisFeatureSink.write()) + .apply(redisFeatureSink.writer()) .getFailedInserts() .apply(Count.globally()); @@ -252,7 +252,7 @@ public void shouldProduceFailedElementIfRetryExceeded() { PCollection failedElementCount = p.apply(Create.of(featureRows)) - .apply(redisFeatureSink.write()) + .apply(redisFeatureSink.writer()) .getFailedInserts() .apply(Count.globally()); @@ -310,7 +310,7 @@ public void shouldConvertRowWithDuplicateEntitiesToValidKey() { .addFields(Field.newBuilder().setValue(Value.newBuilder().setInt64Val(1001))) .build(); - p.apply(Create.of(offendingRow)).apply(redisFeatureSink.write()); + p.apply(Create.of(offendingRow)).apply(redisFeatureSink.writer()); p.run(); @@ -365,7 +365,7 @@ public void shouldConvertRowWithOutOfOrderFieldsToValidKey() { .addAllFields(expectedFields) .build(); - p.apply(Create.of(offendingRow)).apply(redisFeatureSink.write()); + p.apply(Create.of(offendingRow)).apply(redisFeatureSink.writer()); p.run(); @@ -421,7 +421,7 @@ public void shouldMergeDuplicateFeatureFields() { .addFields(Field.newBuilder().setValue(Value.newBuilder().setInt64Val(1001))) .build(); - p.apply(Create.of(featureRowWithDuplicatedFeatureFields)).apply(redisFeatureSink.write()); + p.apply(Create.of(featureRowWithDuplicatedFeatureFields)).apply(redisFeatureSink.writer()); p.run(); @@ -469,7 +469,7 @@ public void shouldPopulateMissingFeatureValuesWithDefaultInstance() { .addFields(Field.newBuilder().setValue(Value.getDefaultInstance())) .build(); - p.apply(Create.of(featureRowWithDuplicatedFeatureFields)).apply(redisFeatureSink.write()); + p.apply(Create.of(featureRowWithDuplicatedFeatureFields)).apply(redisFeatureSink.writer()); p.run(); From 85e3e5718649bbd023a5719224944bd4ff7f70b4 Mon Sep 17 00:00:00 2001 From: zhilingc Date: Wed, 25 Mar 2020 16:35:59 +0800 Subject: [PATCH 2/5] Fix javadoc linting errors --- ingestion/src/test/java/feast/test/TestUtil.java | 3 +++ .../java/feast/storage/api/retrieval/OnlineRetriever.java | 2 +- .../main/java/feast/storage/api/write/WriteResult.java | 8 +++++++- 3 files changed, 11 insertions(+), 2 deletions(-) diff --git a/ingestion/src/test/java/feast/test/TestUtil.java b/ingestion/src/test/java/feast/test/TestUtil.java index 1c9e5bdc555..0a141587078 100644 --- a/ingestion/src/test/java/feast/test/TestUtil.java +++ b/ingestion/src/test/java/feast/test/TestUtil.java @@ -162,6 +162,9 @@ public static void publishFeatureRowsToKafka( * Create a Feature Row with random value according to the FeatureSetSpec * *

See {@link #createRandomFeatureRow(FeatureSetSpec, int)} + * + * @param featureSetSpec {@link FeatureSetSpec} + * @return {@link FeatureRow} */ public static FeatureRow createRandomFeatureRow(FeatureSetSpec featureSetSpec) { ThreadLocalRandom random = ThreadLocalRandom.current(); diff --git a/storage/api/src/main/java/feast/storage/api/retrieval/OnlineRetriever.java b/storage/api/src/main/java/feast/storage/api/retrieval/OnlineRetriever.java index 72888bfc081..591d5c8edcd 100644 --- a/storage/api/src/main/java/feast/storage/api/retrieval/OnlineRetriever.java +++ b/storage/api/src/main/java/feast/storage/api/retrieval/OnlineRetriever.java @@ -32,7 +32,7 @@ public interface OnlineRetriever { * @param entityRows list of entity rows in the feature request * @param featureSetRequests List of {@link FeatureSetRequest} to feature references in the * request tied to that feature set. - * @return list of {@link OnlineRetrieverResponse} for each entity row + * @return list of {@link List} for each entity row */ List> getOnlineFeatures( List entityRows, List featureSetRequests); diff --git a/storage/api/src/main/java/feast/storage/api/write/WriteResult.java b/storage/api/src/main/java/feast/storage/api/write/WriteResult.java index ab92f3600ba..61878ad5850 100644 --- a/storage/api/src/main/java/feast/storage/api/write/WriteResult.java +++ b/storage/api/src/main/java/feast/storage/api/write/WriteResult.java @@ -34,7 +34,13 @@ public final class WriteResult implements Serializable, POutput { private static TupleTag successfulInsertsTag = new TupleTag<>("successfulInserts"); private static TupleTag failedInsertsTupleTag = new TupleTag<>("failedInserts"); - /** Creates a {@link WriteResult} in the given {@link Pipeline}. */ + /** + * /** Creates a {@link WriteResult} in the given {@link Pipeline}. + * @param pipeline {@link Pipeline} + * @param successfulInserts {@link PCollection} of {@link FeatureRow}s successfully inserted into the store + * @param failedInserts {@link PCollection} of {@link FailedElement}s + * @return {@link WriteResult} + */ public static WriteResult in( Pipeline pipeline, PCollection successfulInserts, From 1951bebdccbad449909546abe96b5f77c887ca7a Mon Sep 17 00:00:00 2001 From: zhilingc Date: Wed, 25 Mar 2020 16:43:35 +0800 Subject: [PATCH 3/5] Apply spotless --- .../src/main/java/feast/storage/api/write/WriteResult.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/storage/api/src/main/java/feast/storage/api/write/WriteResult.java b/storage/api/src/main/java/feast/storage/api/write/WriteResult.java index 61878ad5850..8eb9d43c40a 100644 --- a/storage/api/src/main/java/feast/storage/api/write/WriteResult.java +++ b/storage/api/src/main/java/feast/storage/api/write/WriteResult.java @@ -36,8 +36,10 @@ public final class WriteResult implements Serializable, POutput { /** * /** Creates a {@link WriteResult} in the given {@link Pipeline}. + * * @param pipeline {@link Pipeline} - * @param successfulInserts {@link PCollection} of {@link FeatureRow}s successfully inserted into the store + * @param successfulInserts {@link PCollection} of {@link FeatureRow}s successfully inserted into + * the store * @param failedInserts {@link PCollection} of {@link FailedElement}s * @return {@link WriteResult} */ From be837dbfb20fb6ab05ae7156ec97dabc7994a07d Mon Sep 17 00:00:00 2001 From: zhilingc Date: Thu, 26 Mar 2020 10:42:07 +0800 Subject: [PATCH 4/5] Fix javadoc formatting --- .../main/java/feast/storage/api/retrieval/OnlineRetriever.java | 3 ++- .../src/main/java/feast/storage/common/testing/TestUtil.java | 3 ++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/storage/api/src/main/java/feast/storage/api/retrieval/OnlineRetriever.java b/storage/api/src/main/java/feast/storage/api/retrieval/OnlineRetriever.java index 591d5c8edcd..c4c63eb14ca 100644 --- a/storage/api/src/main/java/feast/storage/api/retrieval/OnlineRetriever.java +++ b/storage/api/src/main/java/feast/storage/api/retrieval/OnlineRetriever.java @@ -32,7 +32,8 @@ public interface OnlineRetriever { * @param entityRows list of entity rows in the feature request * @param featureSetRequests List of {@link FeatureSetRequest} to feature references in the * request tied to that feature set. - * @return list of {@link List} for each entity row + * @return list of lists of {@link FeatureRow}s corresponding to each feature set request and + * entity row. */ List> getOnlineFeatures( List entityRows, List featureSetRequests); diff --git a/storage/api/src/main/java/feast/storage/common/testing/TestUtil.java b/storage/api/src/main/java/feast/storage/common/testing/TestUtil.java index d26930fbfed..6047a93dc17 100644 --- a/storage/api/src/main/java/feast/storage/common/testing/TestUtil.java +++ b/storage/api/src/main/java/feast/storage/common/testing/TestUtil.java @@ -33,7 +33,8 @@ public class TestUtil { /** * Create a Feature Row with random value according to the FeatureSetSpec * - *

See {@link #createRandomFeatureRow(FeatureSet, int)} + * @param featureSet {@link FeatureSet} + * @return {@link FeatureRow} */ public static FeatureRow createRandomFeatureRow(FeatureSet featureSet) { ThreadLocalRandom random = ThreadLocalRandom.current(); From d9f6b39730afce2284b71c49c62f9978bfb27a26 Mon Sep 17 00:00:00 2001 From: zhilingc Date: Fri, 27 Mar 2020 15:00:20 +0800 Subject: [PATCH 5/5] Drop result from HistoricalRetrievalResult constructors --- .../storage/api/retrieval/HistoricalRetrievalResult.java | 4 ++-- .../src/main/java/feast/storage/api/write/WriteResult.java | 2 +- .../bigquery/retrieval/BigQueryHistoricalRetriever.java | 6 +++--- 3 files changed, 6 insertions(+), 6 deletions(-) diff --git a/storage/api/src/main/java/feast/storage/api/retrieval/HistoricalRetrievalResult.java b/storage/api/src/main/java/feast/storage/api/retrieval/HistoricalRetrievalResult.java index 59e08dbe8c7..2bf27c6158c 100644 --- a/storage/api/src/main/java/feast/storage/api/retrieval/HistoricalRetrievalResult.java +++ b/storage/api/src/main/java/feast/storage/api/retrieval/HistoricalRetrievalResult.java @@ -48,7 +48,7 @@ public abstract class HistoricalRetrievalResult implements Serializable { * @param error error that occurred * @return {@link HistoricalRetrievalResult} */ - public static HistoricalRetrievalResult errorResult(String id, Exception error) { + public static HistoricalRetrievalResult error(String id, Exception error) { return newBuilder() .setId(id) .setStatus(JobStatus.JOB_STATUS_DONE) @@ -65,7 +65,7 @@ public static HistoricalRetrievalResult errorResult(String id, Exception error) * @param dataFormat data format of the output files * @return */ - public static HistoricalRetrievalResult successResult( + public static HistoricalRetrievalResult success( String id, List fileUris, DataFormat dataFormat) { return newBuilder() .setId(id) diff --git a/storage/api/src/main/java/feast/storage/api/write/WriteResult.java b/storage/api/src/main/java/feast/storage/api/write/WriteResult.java index 8eb9d43c40a..482863ea6e0 100644 --- a/storage/api/src/main/java/feast/storage/api/write/WriteResult.java +++ b/storage/api/src/main/java/feast/storage/api/write/WriteResult.java @@ -35,7 +35,7 @@ public final class WriteResult implements Serializable, POutput { private static TupleTag failedInsertsTupleTag = new TupleTag<>("failedInserts"); /** - * /** Creates a {@link WriteResult} in the given {@link Pipeline}. + * Creates a {@link WriteResult} in the given {@link Pipeline}. * * @param pipeline {@link Pipeline} * @param successfulInserts {@link PCollection} of {@link FeatureRow}s successfully inserted into diff --git a/storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/retrieval/BigQueryHistoricalRetriever.java b/storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/retrieval/BigQueryHistoricalRetriever.java index 46d4fea11c7..63648604331 100644 --- a/storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/retrieval/BigQueryHistoricalRetriever.java +++ b/storage/connectors/bigquery/src/main/java/feast/storage/connectors/bigquery/retrieval/BigQueryHistoricalRetriever.java @@ -100,7 +100,7 @@ public HistoricalRetrievalResult getHistoricalFeatures( TableId entityTableWithUUIDs = generateUUIDs(entityTable); entityTableName = generateFullTableName(entityTableWithUUIDs); } catch (Exception e) { - return HistoricalRetrievalResult.errorResult( + return HistoricalRetrievalResult.error( retrievalId, new RuntimeException( String.format("Unable to load entity table to BigQuery: %s", e.toString()))); @@ -140,12 +140,12 @@ public HistoricalRetrievalResult getHistoricalFeatures( waitForJob(extractJob); } catch (BigQueryException | InterruptedException | IOException e) { - return HistoricalRetrievalResult.errorResult(retrievalId, e); + return HistoricalRetrievalResult.error(retrievalId, e); } List fileUris = parseOutputFileURIs(retrievalId); - return HistoricalRetrievalResult.successResult( + return HistoricalRetrievalResult.success( retrievalId, fileUris, ServingAPIProto.DataFormat.DATA_FORMAT_AVRO); }