From 15b6d1210ac5e07c5b7abdcb5b8810ba2bee2e9e Mon Sep 17 00:00:00 2001 From: David Heryanto Date: Wed, 25 Dec 2019 15:43:22 +0800 Subject: [PATCH 1/6] Update BQ query config to always set destination table, so that it can work with large results Refer to: https://cloud.google.com/bigquery/quotas#query_jobs, maximum reponse-size bullet point. --- .../service/BigQueryServingService.java | 20 ++++++++++- .../bigquery/BatchRetrievalQueryRunnable.java | 36 ++++++++++++++++--- .../store/bigquery/SubqueryCallable.java | 14 ++++++++ 3 files changed, 65 insertions(+), 5 deletions(-) diff --git a/serving/src/main/java/feast/serving/service/BigQueryServingService.java b/serving/src/main/java/feast/serving/service/BigQueryServingService.java index 701e146ee5d..9459f1bbc76 100644 --- a/serving/src/main/java/feast/serving/service/BigQueryServingService.java +++ b/serving/src/main/java/feast/serving/service/BigQueryServingService.java @@ -32,6 +32,7 @@ import com.google.cloud.bigquery.Schema; import com.google.cloud.bigquery.Table; import com.google.cloud.bigquery.TableId; +import com.google.cloud.bigquery.TableInfo; import com.google.cloud.storage.Storage; import feast.core.FeatureSetProto.FeatureSetSpec; import feast.serving.ServingAPIProto; @@ -56,10 +57,13 @@ import java.util.Optional; import java.util.UUID; import java.util.stream.Collectors; +import org.joda.time.Duration; import org.slf4j.Logger; public class BigQueryServingService implements ServingService { + // Default no of millis for which a temporary table should exist before it is deleted in BigQuery. + public static final long TEMP_TABLE_EXPIRY_DURATION_MS = Duration.standardDays(1).getMillis(); private static final Logger log = org.slf4j.LoggerFactory.getLogger(BigQueryServingService.class); private final BigQuery bigquery; @@ -230,9 +234,19 @@ private TableId generateUUIDs(Table loadedEntityTable) { try { String uuidQuery = createEntityTableUUIDQuery(generateFullTableName(loadedEntityTable.getTableId())); - QueryJobConfiguration queryJobConfig = QueryJobConfiguration.newBuilder(uuidQuery).build(); + QueryJobConfiguration queryJobConfig = + QueryJobConfiguration.newBuilder(uuidQuery) + .setDestinationTable(TableId.of(projectId, datasetId, createTempTableName())) + .build(); Job queryJob = bigquery.create(JobInfo.of(queryJobConfig)); queryJob.waitFor(); + TableInfo expiry = + bigquery + .getTable(queryJobConfig.getDestinationTable()) + .toBuilder() + .setExpirationTime(System.currentTimeMillis() + TEMP_TABLE_EXPIRY_DURATION_MS) + .build(); + bigquery.update(expiry); queryJobConfig = queryJob.getConfiguration(); return queryJobConfig.getDestinationTable(); } catch (InterruptedException | BigQueryException e) { @@ -242,4 +256,8 @@ private TableId generateUUIDs(Table loadedEntityTable) { .asRuntimeException(); } } + + public static String createTempTableName() { + return "temp" + UUID.randomUUID().toString().replace("-", ""); + } } diff --git a/serving/src/main/java/feast/serving/store/bigquery/BatchRetrievalQueryRunnable.java b/serving/src/main/java/feast/serving/store/bigquery/BatchRetrievalQueryRunnable.java index 2d51547d0e7..e39b1af07f1 100644 --- a/serving/src/main/java/feast/serving/store/bigquery/BatchRetrievalQueryRunnable.java +++ b/serving/src/main/java/feast/serving/store/bigquery/BatchRetrievalQueryRunnable.java @@ -16,6 +16,8 @@ */ package feast.serving.store.bigquery; +import static feast.serving.service.BigQueryServingService.TEMP_TABLE_EXPIRY_DURATION_MS; +import static feast.serving.service.BigQueryServingService.createTempTableName; import static feast.serving.store.bigquery.QueryTemplater.createTimestampLimitQuery; import com.google.auto.value.AutoValue; @@ -27,6 +29,8 @@ import com.google.cloud.bigquery.Job; import com.google.cloud.bigquery.JobInfo; import com.google.cloud.bigquery.QueryJobConfiguration; +import com.google.cloud.bigquery.TableId; +import com.google.cloud.bigquery.TableInfo; import com.google.cloud.bigquery.TableResult; import com.google.cloud.storage.Blob; import com.google.cloud.storage.Storage; @@ -35,6 +39,7 @@ import feast.serving.ServingAPIProto.DataFormat; import feast.serving.ServingAPIProto.JobStatus; import feast.serving.ServingAPIProto.JobType; +import feast.serving.service.BigQueryServingService; import feast.serving.service.JobService; import feast.serving.store.bigquery.model.FeatureSetInfo; import io.grpc.Status; @@ -175,15 +180,17 @@ Job runBatchQuery(List featureSetQueries) ExecutorCompletionService executorCompletionService = new ExecutorCompletionService<>(executorService); - List featureSetInfos = new ArrayList<>(); for (int i = 0; i < featureSetQueries.size(); i++) { QueryJobConfiguration queryJobConfig = - QueryJobConfiguration.newBuilder(featureSetQueries.get(i)).build(); + QueryJobConfiguration.newBuilder(featureSetQueries.get(i)) + .setDestinationTable(TableId.of(projectId(), datasetId(), createTempTableName())) + .build(); Job subqueryJob = bigquery().create(JobInfo.of(queryJobConfig)); executorCompletionService.submit( SubqueryCallable.builder() + .setBigquery(bigquery()) .setFeatureSetInfo(featureSetInfos().get(i)) .setSubqueryJob(subqueryJob) .build()); @@ -191,7 +198,8 @@ Job runBatchQuery(List featureSetQueries) for (int i = 0; i < featureSetQueries.size(); i++) { try { - FeatureSetInfo featureSetInfo = executorCompletionService.take().get(SUBQUERY_TIMEOUT_SECS, TimeUnit.SECONDS); + FeatureSetInfo featureSetInfo = + executorCompletionService.take().get(SUBQUERY_TIMEOUT_SECS, TimeUnit.SECONDS); featureSetInfos.add(featureSetInfo); } catch (InterruptedException | ExecutionException | TimeoutException e) { jobService() @@ -214,9 +222,20 @@ Job runBatchQuery(List featureSetQueries) String joinQuery = QueryTemplater.createJoinQuery( featureSetInfos, entityTableColumnNames(), entityTableName()); - QueryJobConfiguration queryJobConfig = QueryJobConfiguration.newBuilder(joinQuery).build(); + QueryJobConfiguration queryJobConfig = + QueryJobConfiguration.newBuilder(joinQuery) + .setDestinationTable(TableId.of(projectId(), datasetId(), createTempTableName())) + .build(); queryJob = bigquery().create(JobInfo.of(queryJobConfig)); queryJob.waitFor(); + TableInfo expiry = + bigquery() + .getTable(queryJobConfig.getDestinationTable()) + .toBuilder() + .setExpirationTime( + System.currentTimeMillis() + BigQueryServingService.TEMP_TABLE_EXPIRY_DURATION_MS) + .build(); + bigquery().update(expiry); return queryJob; } @@ -248,10 +267,19 @@ private FieldValueList getTimestampLimits(String entityTableName) { QueryJobConfiguration getTimestampLimitsQuery = QueryJobConfiguration.newBuilder(createTimestampLimitQuery(entityTableName)) .setDefaultDataset(DatasetId.of(projectId(), datasetId())) + .setDestinationTable(TableId.of(projectId(), datasetId(), createTempTableName())) .build(); try { Job job = bigquery().create(JobInfo.of(getTimestampLimitsQuery)); TableResult getTimestampLimitsQueryResult = job.waitFor().getQueryResults(); + TableInfo expiry = + bigquery() + .getTable(getTimestampLimitsQuery.getDestinationTable()) + .toBuilder() + .setExpirationTime(System.currentTimeMillis() + TEMP_TABLE_EXPIRY_DURATION_MS) + .build(); + bigquery().update(expiry); + FieldValueList result = null; for (FieldValueList fields : getTimestampLimitsQueryResult.getValues()) { result = fields; diff --git a/serving/src/main/java/feast/serving/store/bigquery/SubqueryCallable.java b/serving/src/main/java/feast/serving/store/bigquery/SubqueryCallable.java index 3c28194e7a3..b2a9009a749 100644 --- a/serving/src/main/java/feast/serving/store/bigquery/SubqueryCallable.java +++ b/serving/src/main/java/feast/serving/store/bigquery/SubqueryCallable.java @@ -16,13 +16,16 @@ */ package feast.serving.store.bigquery; +import static feast.serving.service.BigQueryServingService.TEMP_TABLE_EXPIRY_DURATION_MS; import static feast.serving.store.bigquery.QueryTemplater.generateFullTableName; import com.google.auto.value.AutoValue; +import com.google.cloud.bigquery.BigQuery; import com.google.cloud.bigquery.BigQueryException; import com.google.cloud.bigquery.Job; import com.google.cloud.bigquery.QueryJobConfiguration; import com.google.cloud.bigquery.TableId; +import com.google.cloud.bigquery.TableInfo; import feast.serving.store.bigquery.model.FeatureSetInfo; import java.util.concurrent.Callable; @@ -33,6 +36,8 @@ @AutoValue public abstract class SubqueryCallable implements Callable { + public abstract BigQuery bigquery(); + public abstract FeatureSetInfo featureSetInfo(); public abstract Job subqueryJob(); @@ -44,6 +49,8 @@ public static Builder builder() { @AutoValue.Builder public abstract static class Builder { + public abstract Builder setBigquery(BigQuery bigquery); + public abstract Builder setFeatureSetInfo(FeatureSetInfo featureSetInfo); public abstract Builder setSubqueryJob(Job subqueryJob); @@ -57,6 +64,13 @@ public FeatureSetInfo call() throws BigQueryException, InterruptedException { subqueryJob().waitFor(); subqueryConfig = subqueryJob().getConfiguration(); TableId destinationTable = subqueryConfig.getDestinationTable(); + TableInfo expiry = + bigquery() + .getTable(destinationTable) + .toBuilder() + .setExpirationTime(System.currentTimeMillis() + TEMP_TABLE_EXPIRY_DURATION_MS) + .build(); + bigquery().update(expiry); String fullTablePath = generateFullTableName(destinationTable); return new FeatureSetInfo(featureSetInfo(), fullTablePath); From 1360d083f14b125c0f265aeff8b5e0c7fd8bef0a Mon Sep 17 00:00:00 2001 From: David Heryanto Date: Wed, 25 Dec 2019 16:03:19 +0800 Subject: [PATCH 2/6] Replace prefix for temp table name --- .../main/java/feast/serving/service/BigQueryServingService.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/serving/src/main/java/feast/serving/service/BigQueryServingService.java b/serving/src/main/java/feast/serving/service/BigQueryServingService.java index 9459f1bbc76..26775e1c541 100644 --- a/serving/src/main/java/feast/serving/service/BigQueryServingService.java +++ b/serving/src/main/java/feast/serving/service/BigQueryServingService.java @@ -258,6 +258,6 @@ private TableId generateUUIDs(Table loadedEntityTable) { } public static String createTempTableName() { - return "temp" + UUID.randomUUID().toString().replace("-", ""); + return "_" + UUID.randomUUID().toString().replace("-", ""); } } From 937a26566311758d56326ad047abae9426f011a6 Mon Sep 17 00:00:00 2001 From: David Heryanto Date: Wed, 25 Dec 2019 19:45:31 +0800 Subject: [PATCH 3/6] Set expiry on entity rows table --- .../service/BigQueryServingService.java | 24 +++++++++---------- 1 file changed, 12 insertions(+), 12 deletions(-) diff --git a/serving/src/main/java/feast/serving/service/BigQueryServingService.java b/serving/src/main/java/feast/serving/service/BigQueryServingService.java index 26775e1c541..81b52ed6538 100644 --- a/serving/src/main/java/feast/serving/service/BigQueryServingService.java +++ b/serving/src/main/java/feast/serving/service/BigQueryServingService.java @@ -186,15 +186,15 @@ private Table loadEntities(DatasetSource datasetSource) { switch (datasetSource.getDatasetSourceCase()) { case FILE_SOURCE: try { - String tableName = generateTemporaryTableName(); - log.info("Loading entity dataset to table {}.{}.{}", projectId, datasetId, tableName); - TableId tableId = TableId.of(projectId, datasetId, tableName); - // Currently only avro supported + // Currently only AVRO format is supported if (datasetSource.getFileSource().getDataFormat() != DataFormat.DATA_FORMAT_AVRO) { throw Status.INVALID_ARGUMENT - .withDescription("Invalid file format, only avro supported") + .withDescription("Invalid file format, only AVRO is supported.") .asRuntimeException(); } + + TableId tableId = TableId.of(projectId, datasetId, createTempTableName()); + log.info("Loading entity rows to: {}.{}.{}", projectId, datasetId, tableId.getTable()); LoadJobConfiguration loadJobConfiguration = LoadJobConfiguration.of( tableId, datasetSource.getFileSource().getFileUrisList(), FormatOptions.avro()); @@ -202,6 +202,13 @@ private Table loadEntities(DatasetSource datasetSource) { loadJobConfiguration.toBuilder().setUseAvroLogicalTypes(true).build(); Job job = bigquery.create(JobInfo.of(loadJobConfiguration)); job.waitFor(); + TableInfo expiry = + bigquery + .getTable(tableId) + .toBuilder() + .setExpirationTime(System.currentTimeMillis() + TEMP_TABLE_EXPIRY_DURATION_MS) + .build(); + bigquery.update(expiry); loadedEntityTable = bigquery.getTable(tableId); if (!loadedEntityTable.exists()) { throw new RuntimeException( @@ -223,13 +230,6 @@ private Table loadEntities(DatasetSource datasetSource) { } } - private String generateTemporaryTableName() { - String source = String.format("feastserving%d", System.currentTimeMillis()); - String guid = UUID.nameUUIDFromBytes(source.getBytes()).toString(); - String suffix = guid.substring(0, Math.min(guid.length(), 10)).replaceAll("-", ""); - return String.format("temp_%s", suffix); - } - private TableId generateUUIDs(Table loadedEntityTable) { try { String uuidQuery = From 478407ada039bdf405b5ccbdf0fd4a140a8d2931 Mon Sep 17 00:00:00 2001 From: David Heryanto Date: Wed, 25 Dec 2019 19:58:42 +0800 Subject: [PATCH 4/6] Include exception message in error description GRPC client such as Feast Python SDK will usually not show error cause only error description --- .../main/java/feast/serving/service/BigQueryServingService.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/serving/src/main/java/feast/serving/service/BigQueryServingService.java b/serving/src/main/java/feast/serving/service/BigQueryServingService.java index 81b52ed6538..7a950e3c8a9 100644 --- a/serving/src/main/java/feast/serving/service/BigQueryServingService.java +++ b/serving/src/main/java/feast/serving/service/BigQueryServingService.java @@ -218,7 +218,7 @@ private Table loadEntities(DatasetSource datasetSource) { } catch (Exception e) { log.error("Exception has occurred in loadEntities method: ", e); throw Status.INTERNAL - .withDescription("Failed to load entity dataset into store") + .withDescription("Failed to load entity dataset into store: " + e.toString()) .withCause(e) .asRuntimeException(); } From c52c549480b4aa62ac458e004a0b12cf6238f2fe Mon Sep 17 00:00:00 2001 From: David Heryanto Date: Wed, 25 Dec 2019 20:14:22 +0800 Subject: [PATCH 5/6] Code cleanup --- .../serving/store/bigquery/BatchRetrievalQueryRunnable.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/serving/src/main/java/feast/serving/store/bigquery/BatchRetrievalQueryRunnable.java b/serving/src/main/java/feast/serving/store/bigquery/BatchRetrievalQueryRunnable.java index e39b1af07f1..47587e1d0ef 100644 --- a/serving/src/main/java/feast/serving/store/bigquery/BatchRetrievalQueryRunnable.java +++ b/serving/src/main/java/feast/serving/store/bigquery/BatchRetrievalQueryRunnable.java @@ -39,7 +39,6 @@ import feast.serving.ServingAPIProto.DataFormat; import feast.serving.ServingAPIProto.JobStatus; import feast.serving.ServingAPIProto.JobType; -import feast.serving.service.BigQueryServingService; import feast.serving.service.JobService; import feast.serving.store.bigquery.model.FeatureSetInfo; import io.grpc.Status; @@ -233,7 +232,7 @@ Job runBatchQuery(List featureSetQueries) .getTable(queryJobConfig.getDestinationTable()) .toBuilder() .setExpirationTime( - System.currentTimeMillis() + BigQueryServingService.TEMP_TABLE_EXPIRY_DURATION_MS) + System.currentTimeMillis() + TEMP_TABLE_EXPIRY_DURATION_MS) .build(); bigquery().update(expiry); From 5381872ae7a3092561ce002c4d151f7f482de127 Mon Sep 17 00:00:00 2001 From: David Heryanto Date: Wed, 25 Dec 2019 21:59:46 +0800 Subject: [PATCH 6/6] Update batch-retrieval e2e test. Output rows may not have the same order as requested entity rows --- tests/e2e/bq-batch-retrieval.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/tests/e2e/bq-batch-retrieval.py b/tests/e2e/bq-batch-retrieval.py index 067dd14a2fb..2d6668eaa86 100644 --- a/tests/e2e/bq-batch-retrieval.py +++ b/tests/e2e/bq-batch-retrieval.py @@ -14,6 +14,7 @@ from feast.type_map import ValueType from google.protobuf.duration_pb2 import Duration +pd.set_option('display.max_columns', None) @pytest.fixture(scope="module") def core_url(pytestconfig): @@ -112,8 +113,8 @@ def test_additional_columns_in_entity_table(client): feature_retrieval_job = client.get_batch_features( entity_rows=entity_df, feature_ids=["additional_columns:1:feature_value"] ) - output = feature_retrieval_job.to_dataframe() - print(output.head()) + output = feature_retrieval_job.to_dataframe().sort_values(by=["entity_id"]) + print(output.head(10)) assert np.allclose(output["additional_float_col"], entity_df["additional_float_col"]) assert output["additional_string_col"].to_list() == entity_df["additional_string_col"].to_list()