This is an automated email from the ASF dual-hosted git repository.
Abacn pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 6b677a58ee9 Update list of stale bq dataset (#40286)
6b677a58ee9 is described below
commit 6b677a58ee94dc72a5dae0bf4725698299672c0d
Author: Yi Hu <[email protected]>
AuthorDate: Wed Sep 30 11:57:03 2026 -0400
Update list of stale bq dataset (#40286)
* Update list of stale bq dataset
* Fix dataset leak in `BigqueryClient.deleteDataset` when a test dataset
contains >50 tables (e.g. `StorageApiSinkCreateIfNeededIT`
due to response page truncation
* Use a unified prefix for integration tests to avoid temp dataset
---
.test-infra/tools/stale_bq_datasets_cleaner.sh | 2 +-
.../apache/beam/examples/complete/TrafficMaxLaneFlowIT.java | 3 ++-
.../apache/beam/examples/cookbook/MaxPerKeyExamplesIT.java | 3 ++-
.../transforms/io/gcp/bigquery/BigQuerySamplesIT.java | 7 ++++++-
.../apache/beam/it/gcp/bigquery/BigQueryStreamingLT.java | 3 ++-
.../sql/meta/provider/iceberg/IcebergReadWriteIT.java | 3 ++-
.../sql/meta/provider/iceberg/PubsubToIcebergIT.java | 3 ++-
.../org/apache/beam/sdk/io/gcp/testing/BigqueryClient.java | 13 +------------
.../org/apache/beam/sdk/io/gcp/testing/BigtableUtils.java | 2 ++
.../apache/beam/sdk/io/gcp/bigquery/BigQueryKmsKeyIT.java | 7 ++++++-
.../sdk/io/gcp/bigquery/BigQuerySchemaUpdateOptionsIT.java | 4 +++-
.../gcp/bigquery/BigQueryTimePartitioningClusteringIT.java | 5 ++++-
.../beam/sdk/io/gcp/bigquery/BigQueryTimestampPicosIT.java | 8 +++++++-
.../apache/beam/sdk/io/gcp/bigquery/BigQueryToTableIT.java | 7 ++++++-
.../gcp/bigquery/StorageApiDataTriggeredSchemaUpdateIT.java | 3 ++-
.../sdk/io/gcp/bigquery/StorageApiSinkCreateIfNeededIT.java | 3 ++-
.../sdk/io/gcp/bigquery/StorageApiSinkFailedRowsIT.java | 3 ++-
.../StorageApiSinkSchemaUpdateWithInputSchemaIT.java | 4 +++-
.../StorageApiSinkSchemaUpdateWithoutInputSchemaIT.java | 4 +++-
.../sdk/io/gcp/bigquery/providers/BigQueryManagedIT.java | 4 +++-
.../apache_beam/io/gcp/bigquery_change_history_it_test.py | 4 ++--
sdks/python/apache_beam/io/gcp/bigquery_file_loads_test.py | 5 +++--
sdks/python/apache_beam/io/gcp/bigquery_read_it_test.py | 11 ++++++-----
sdks/python/apache_beam/io/gcp/bigquery_test.py | 6 +++---
sdks/python/apache_beam/io/gcp/bigquery_tools.py | 3 +++
.../ml/inference/vertex_ai_model_monitoring_v2_it_test.py | 5 +++--
26 files changed, 81 insertions(+), 44 deletions(-)
diff --git a/.test-infra/tools/stale_bq_datasets_cleaner.sh
b/.test-infra/tools/stale_bq_datasets_cleaner.sh
index 652deaddcc5..36f77239741 100755
--- a/.test-infra/tools/stale_bq_datasets_cleaner.sh
+++ b/.test-infra/tools/stale_bq_datasets_cleaner.sh
@@ -24,7 +24,7 @@ PROJECT=apache-beam-testing
MAX_RESULT=1500
BQ_DATASETS=`bq --project_id=$PROJECT ls --max_results=$MAX_RESULT | tail -n
$MAX_RESULT | sed s/^[[:space:]]*/${PROJECT}:/`
-CLEANUP_DATASET_TEMPLATES=(beam_bigquery_samples_ beam_temp_dataset_
FHIR_store_ bq_query_schema_update_options_16 bq_query_to_table_16
'\:(bq_read_all_|combine_per_key_examples|filter_examples|hourly_team_score_|leader_board_|leaderboard_|game_stats_|python_|temp_dataset)[a-z_]*[0-9a-f]{12,}$')
+CLEANUP_DATASET_TEMPLATES=( beam_temp_dataset_ file_loads_streaming_it_
FHIR_store_ io_bigquery_lt_ managed-iceberg-biglake-its.biglake_test_catalog_
'managed_iceberg_bqms.*_tests?_' '_(python|java)_validations_'
'\:(bq_read_all_|combine_per_key_examples|filter_examples|hourly_team_score_|leader_board_|leaderboard_|game_stats_|python_|temp_dataset)[a-z_]*[0-9a-f]{12,}$')
# A grace period of 5 days
GRACE_PERIOD=$((`date +%s` - 24 * 3600 * 5))
diff --git
a/examples/java/src/test/java/org/apache/beam/examples/complete/TrafficMaxLaneFlowIT.java
b/examples/java/src/test/java/org/apache/beam/examples/complete/TrafficMaxLaneFlowIT.java
index 1832d00aec1..402ef13e944 100644
---
a/examples/java/src/test/java/org/apache/beam/examples/complete/TrafficMaxLaneFlowIT.java
+++
b/examples/java/src/test/java/org/apache/beam/examples/complete/TrafficMaxLaneFlowIT.java
@@ -17,6 +17,7 @@
*/
package org.apache.beam.examples.complete;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
import static org.junit.Assert.assertEquals;
import com.google.api.client.util.BackOff;
@@ -44,7 +45,7 @@ public class TrafficMaxLaneFlowIT {
private TrafficMaxLaneFlowOptions options;
private final String timestamp = Long.toString(System.currentTimeMillis());
- private final String outputDatasetId = "traffic_max_lane_flow_" + timestamp;
+ private final String outputDatasetId = TEMP_DATASET_PREFIX +
"traffic_max_lane_flow_" + timestamp;
private final String outputTable = "traffic_max_lane_flow_table";
private String projectId;
private BigqueryClient bqClient;
diff --git
a/examples/java/src/test/java/org/apache/beam/examples/cookbook/MaxPerKeyExamplesIT.java
b/examples/java/src/test/java/org/apache/beam/examples/cookbook/MaxPerKeyExamplesIT.java
index 16ea27d9132..53eb3413d57 100644
---
a/examples/java/src/test/java/org/apache/beam/examples/cookbook/MaxPerKeyExamplesIT.java
+++
b/examples/java/src/test/java/org/apache/beam/examples/cookbook/MaxPerKeyExamplesIT.java
@@ -17,6 +17,7 @@
*/
package org.apache.beam.examples.cookbook;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
import static org.junit.Assert.assertEquals;
import com.google.api.client.util.BackOff;
@@ -42,7 +43,7 @@ import org.junit.runners.JUnit4;
public class MaxPerKeyExamplesIT {
private MaxPerKeyExamplesIT.MaxPerKeyExamplesOptions options;
private final String timestamp = Long.toString(System.currentTimeMillis());
- private final String outputDatasetId = "max_per_key_examples" + timestamp;
+ private final String outputDatasetId = TEMP_DATASET_PREFIX + "maxperkey_" +
timestamp;
private final String outputTable = "max_per_key_examples_table";
private final Long defaultExpiration = 1000L * 60 * 60;
private String projectId;
diff --git
a/examples/java/src/test/java/org/apache/beam/examples/snippets/transforms/io/gcp/bigquery/BigQuerySamplesIT.java
b/examples/java/src/test/java/org/apache/beam/examples/snippets/transforms/io/gcp/bigquery/BigQuerySamplesIT.java
index cdaffa901bc..a1434236277 100644
---
a/examples/java/src/test/java/org/apache/beam/examples/snippets/transforms/io/gcp/bigquery/BigQuerySamplesIT.java
+++
b/examples/java/src/test/java/org/apache/beam/examples/snippets/transforms/io/gcp/bigquery/BigQuerySamplesIT.java
@@ -17,6 +17,7 @@
*/
package org.apache.beam.examples.snippets.transforms.io.gcp.bigquery;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
import static org.junit.Assert.assertEquals;
import com.google.api.services.bigquery.model.TableRow;
@@ -76,7 +77,11 @@ public class BigQuerySamplesIT {
private static final BigQuery BIGQUERY =
BigQueryOptions.newBuilder().setProjectId(PROJECT).build().getService();
private static final String DATASET =
- "beam_bigquery_samples_" + System.currentTimeMillis() + "_" + new
SecureRandom().nextInt(32);
+ TEMP_DATASET_PREFIX
+ + "samples_"
+ + System.currentTimeMillis()
+ + "_"
+ + new SecureRandom().nextInt(32);
@Rule public final transient TestPipeline writePipeline =
TestPipeline.create();
@Rule public final transient TestPipeline readTablePipeline =
TestPipeline.create();
diff --git
a/it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/bigquery/BigQueryStreamingLT.java
b/it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/bigquery/BigQueryStreamingLT.java
index 1b7257a8cde..7715f782b54 100644
---
a/it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/bigquery/BigQueryStreamingLT.java
+++
b/it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/bigquery/BigQueryStreamingLT.java
@@ -19,6 +19,7 @@ package org.apache.beam.it.gcp.bigquery;
import static
org.apache.beam.sdk.io.gcp.bigquery.BigQueryUtils.toTableReference;
import static org.apache.beam.sdk.io.gcp.bigquery.BigQueryUtils.toTableSpec;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotEquals;
@@ -95,7 +96,7 @@ public class BigQueryStreamingLT extends IOLoadTestBase {
private static final BigqueryClient BQ_CLIENT = new
BigqueryClient("BigQueryStreamingLT");
private static final String BIG_QUERY_DATASET_ID =
- "storage_api_sink_load_test_" + System.nanoTime();
+ TEMP_DATASET_PREFIX + "sink_load_test_" + System.nanoTime();
private TestConfiguration config;
private Integer crashIntervalSeconds;
diff --git
a/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergReadWriteIT.java
b/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergReadWriteIT.java
index 3a791c6fe88..5b47fca79ce 100644
---
a/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergReadWriteIT.java
+++
b/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/IcebergReadWriteIT.java
@@ -19,6 +19,7 @@ package
org.apache.beam.sdk.extensions.sql.meta.provider.iceberg;
import static java.lang.String.format;
import static java.util.Arrays.asList;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
import static org.apache.beam.sdk.schemas.Schema.FieldType.BOOLEAN;
import static org.apache.beam.sdk.schemas.Schema.FieldType.DOUBLE;
import static org.apache.beam.sdk.schemas.Schema.FieldType.FLOAT;
@@ -101,7 +102,7 @@ public class IcebergReadWriteIT {
private static final BigqueryClient BQ_CLIENT = new
BigqueryClient("IcebergReadWriteIT");
private static final String BQMS_CATALOG =
"org.apache.iceberg.gcp.bigquery.BigQueryMetastoreCatalog";
- static final String DATASET = "iceberg_sql_tests_" + System.nanoTime();
+ static final String DATASET = TEMP_DATASET_PREFIX + "iceberg_sql_" +
System.nanoTime();
static String warehouse;
protected static final GcpOptions OPTIONS =
TestPipeline.testingPipelineOptions().as(GcpOptions.class);
diff --git
a/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/PubsubToIcebergIT.java
b/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/PubsubToIcebergIT.java
index 8b250af2754..5d9475a8c4b 100644
---
a/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/PubsubToIcebergIT.java
+++
b/sdks/java/extensions/sql/iceberg/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/iceberg/PubsubToIcebergIT.java
@@ -19,6 +19,7 @@ package
org.apache.beam.sdk.extensions.sql.meta.provider.iceberg;
import static java.lang.String.format;
import static java.nio.charset.StandardCharsets.UTF_8;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
import static org.apache.beam.sdk.schemas.Schema.FieldType.INT64;
import static org.apache.beam.sdk.schemas.Schema.FieldType.STRING;
import static org.junit.Assert.assertEquals;
@@ -76,7 +77,7 @@ public class PubsubToIcebergIT implements Serializable {
private static final BigqueryClient BQ_CLIENT = new
BigqueryClient("PubsubToIcebergIT");
private static final String BQMS_CATALOG =
"org.apache.iceberg.gcp.bigquery.BigQueryMetastoreCatalog";
- static final String DATASET = "sql_pubsub_to_iceberg_it_" +
System.nanoTime();
+ static final String DATASET = TEMP_DATASET_PREFIX + "sql_pubsub_to_iceberg_"
+ System.nanoTime();
static String warehouse;
private static Catalog icebergCatalog;
protected static final GcpOptions OPTIONS =
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/testing/BigqueryClient.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/testing/BigqueryClient.java
index 151d94f0349..a83c4940491 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/testing/BigqueryClient.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/testing/BigqueryClient.java
@@ -43,8 +43,6 @@ import
com.google.api.services.bigquery.model.TableDataInsertAllRequest;
import com.google.api.services.bigquery.model.TableDataInsertAllRequest.Rows;
import com.google.api.services.bigquery.model.TableDataInsertAllResponse;
import com.google.api.services.bigquery.model.TableFieldSchema;
-import com.google.api.services.bigquery.model.TableList;
-import com.google.api.services.bigquery.model.TableList.Tables;
import com.google.api.services.bigquery.model.TableReference;
import com.google.api.services.bigquery.model.TableRow;
import com.google.api.services.bigquery.model.TableSchema;
@@ -482,16 +480,7 @@ public class BigqueryClient {
public void deleteDataset(String projectId, String datasetId) {
try {
- TableList tables = bqClient.tables().list(projectId,
datasetId).execute();
- for (Tables table : tables.getTables()) {
- this.deleteTable(projectId, datasetId,
table.getTableReference().getTableId());
- }
- } catch (Exception e) {
- LOG.debug("Exceptions caught when listing all tables", e);
- }
-
- try {
- bqClient.datasets().delete(projectId, datasetId).execute();
+ bqClient.datasets().delete(projectId,
datasetId).setDeleteContents(true).execute();
LOG.info("Successfully deleted dataset: {}", datasetId);
} catch (Exception e) {
LOG.debug("Exceptions caught when deleting dataset", e);
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/testing/BigtableUtils.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/testing/BigtableUtils.java
index 0d1a5e05b17..3cf45cfc001 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/testing/BigtableUtils.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/testing/BigtableUtils.java
@@ -23,6 +23,8 @@ import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.primitives.Lon
public class BigtableUtils {
+ public static final String TEMP_DATASET_PREFIX = "beam_temp_dataset_";
+
public static ByteString byteString(byte[] bytes) {
return ByteString.copyFrom(bytes);
}
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryKmsKeyIT.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryKmsKeyIT.java
index edf93dba3ae..b85b6d1e54c 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryKmsKeyIT.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryKmsKeyIT.java
@@ -17,6 +17,7 @@
*/
package org.apache.beam.sdk.io.gcp.bigquery;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
@@ -56,7 +57,11 @@ public class BigQueryKmsKeyIT {
private static final BigqueryClient BQ_CLIENT = new
BigqueryClient("BigQueryKmsKeyIT");
private static final String BIG_QUERY_DATASET_ID =
- "bq_query_to_table_" + System.currentTimeMillis() + "_" + new
SecureRandom().nextInt(32);
+ TEMP_DATASET_PREFIX
+ + "query_to_table_"
+ + System.currentTimeMillis()
+ + "_"
+ + new SecureRandom().nextInt(32);
private static final TableSchema OUTPUT_SCHEMA =
new TableSchema()
.setFields(ImmutableList.of(new
TableFieldSchema().setName("fruit").setType("STRING")));
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQuerySchemaUpdateOptionsIT.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQuerySchemaUpdateOptionsIT.java
index c5e954a4638..9ac2cc4bc08 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQuerySchemaUpdateOptionsIT.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQuerySchemaUpdateOptionsIT.java
@@ -17,6 +17,7 @@
*/
package org.apache.beam.sdk.io.gcp.bigquery;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
import static org.junit.Assert.assertEquals;
import com.google.api.services.bigquery.model.QueryResponse;
@@ -68,7 +69,8 @@ public class BigQuerySchemaUpdateOptionsIT {
new BigqueryClient("BigQuerySchemaUpdateOptionsIT");
private static final String BIG_QUERY_DATASET_ID =
- "bq_query_schema_update_options_"
+ TEMP_DATASET_PREFIX
+ + "schema_update_options_"
+ System.currentTimeMillis()
+ "_"
+ new SecureRandom().nextInt(32);
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryTimePartitioningClusteringIT.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryTimePartitioningClusteringIT.java
index 008cc02beee..848d4c2ea6e 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryTimePartitioningClusteringIT.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryTimePartitioningClusteringIT.java
@@ -17,6 +17,8 @@
*/
package org.apache.beam.sdk.io.gcp.bigquery;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
+
import com.google.api.services.bigquery.Bigquery;
import com.google.api.services.bigquery.model.Clustering;
import com.google.api.services.bigquery.model.Table;
@@ -59,7 +61,8 @@ public class BigQueryTimePartitioningClusteringIT {
private static final BigqueryClient BQ_CLIENT =
new BigqueryClient("BigQueryTimePartitioningClusteringIT");
private static final String DATASET_NAME =
- "BigQueryTimePartitioningIT_"
+ TEMP_DATASET_PREFIX
+ + "time_partitioning_"
+ System.currentTimeMillis()
+ "_"
+ new SecureRandom().nextInt(32);
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryTimestampPicosIT.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryTimestampPicosIT.java
index 07b6adf46bc..a5303240052 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryTimestampPicosIT.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryTimestampPicosIT.java
@@ -17,6 +17,8 @@
*/
package org.apache.beam.sdk.io.gcp.bigquery;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
+
import com.google.api.services.bigquery.model.TableFieldSchema;
import com.google.api.services.bigquery.model.TableRow;
import com.google.api.services.bigquery.model.TableSchema;
@@ -56,7 +58,11 @@ public class BigQueryTimestampPicosIT {
private static String project;
private static final String DATASET_ID =
- "bq_ts_picos_" + System.currentTimeMillis() + "_" + new
SecureRandom().nextInt(32);
+ TEMP_DATASET_PREFIX
+ + "ts_picos_"
+ + System.currentTimeMillis()
+ + "_"
+ + new SecureRandom().nextInt(32);
private static final BigqueryClient BQ_CLIENT = new
BigqueryClient("BigQueryTimestampPicosIT");
private static TestBigQueryOptions bqOptions;
private static String nestedTableSpec;
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryToTableIT.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryToTableIT.java
index 31bb627e845..407f00142cd 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryToTableIT.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryToTableIT.java
@@ -17,6 +17,7 @@
*/
package org.apache.beam.sdk.io.gcp.bigquery;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
import static org.junit.Assert.assertEquals;
import com.google.api.client.util.BackOff;
@@ -72,7 +73,11 @@ public class BigQueryToTableIT {
private static final BigqueryClient BQ_CLIENT = new
BigqueryClient("BigQueryToTableIT");
private static final String BIG_QUERY_DATASET_ID =
- "bq_query_to_table_" + System.currentTimeMillis() + "_" + new
SecureRandom().nextInt(32);
+ TEMP_DATASET_PREFIX
+ + "query_to_table_"
+ + System.currentTimeMillis()
+ + "_"
+ + new SecureRandom().nextInt(32);
private static final TableSchema LEGACY_QUERY_TABLE_SCHEMA =
new TableSchema()
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiDataTriggeredSchemaUpdateIT.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiDataTriggeredSchemaUpdateIT.java
index d90dedaf06f..25cc017a1ca 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiDataTriggeredSchemaUpdateIT.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiDataTriggeredSchemaUpdateIT.java
@@ -17,6 +17,7 @@
*/
package org.apache.beam.sdk.io.gcp.bigquery;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
import static
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects.firstNonNull;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertTrue;
@@ -82,7 +83,7 @@ public class StorageApiDataTriggeredSchemaUpdateIT {
private static final String PROJECT =
TestPipeline.testingPipelineOptions().as(GcpOptions.class).getProject();
private static final String BIG_QUERY_DATASET_ID =
- "storage_api_data_triggered_schema_update_" + System.nanoTime();
+ TEMP_DATASET_PREFIX + "triggered_schema_update_" + System.nanoTime();
private static String bigQueryLocation;
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkCreateIfNeededIT.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkCreateIfNeededIT.java
index 858921e19ce..70eae409598 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkCreateIfNeededIT.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkCreateIfNeededIT.java
@@ -19,6 +19,7 @@ package org.apache.beam.sdk.io.gcp.bigquery;
import static org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO.CONNECTION_ID;
import static org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO.STORAGE_URI;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.junit.Assert.assertEquals;
@@ -71,7 +72,7 @@ public class StorageApiSinkCreateIfNeededIT {
private static final String PROJECT =
TestPipeline.testingPipelineOptions().as(GcpOptions.class).getProject();
private static final String BIG_QUERY_DATASET_ID =
- "storage_api_sink_create_tables_" + System.nanoTime();
+ TEMP_DATASET_PREFIX + "storageapi_" + System.nanoTime();
private static final String TEST_CONNECTION_ID =
"projects/apache-beam-testing/locations/us/connections/apache-beam-testing-storageapi-biglake-nodelete";
private static final String TEST_STORAGE_URI =
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkFailedRowsIT.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkFailedRowsIT.java
index f721f57147e..5124e46847d 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkFailedRowsIT.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkFailedRowsIT.java
@@ -17,6 +17,7 @@
*/
package org.apache.beam.sdk.io.gcp.bigquery;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
import static org.hamcrest.MatcherAssert.assertThat;
import com.google.api.services.bigquery.model.Table;
@@ -74,7 +75,7 @@ public class StorageApiSinkFailedRowsIT {
private static final String PROJECT =
TestPipeline.testingPipelineOptions().as(GcpOptions.class).getProject();
private static final String BIG_QUERY_DATASET_ID =
- "storage_api_sink_failed_rows" + System.nanoTime();
+ TEMP_DATASET_PREFIX + "failed_rows" + System.nanoTime();
private static final List<TableFieldSchema> FIELDS =
ImmutableList.<TableFieldSchema>builder()
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkSchemaUpdateWithInputSchemaIT.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkSchemaUpdateWithInputSchemaIT.java
index 169fb4afc6e..6bad805f045 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkSchemaUpdateWithInputSchemaIT.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkSchemaUpdateWithInputSchemaIT.java
@@ -17,6 +17,8 @@
*/
package org.apache.beam.sdk.io.gcp.bigquery;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
+
import java.io.IOException;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
import org.junit.AfterClass;
@@ -30,7 +32,7 @@ import org.junit.runners.Parameterized;
@RunWith(StorageApiSinkSchemaUpdateITBase.ParallelParameterized.class)
public class StorageApiSinkSchemaUpdateWithInputSchemaIT extends
StorageApiSinkSchemaUpdateITBase {
private static final String BIG_QUERY_DATASET_ID =
- "storage_api_sink_schema_change_with_input_" + System.nanoTime();
+ TEMP_DATASET_PREFIX + "sink_schema_change_with_input_" +
System.nanoTime();
@Parameterized.Parameters(name = "changeTableSchema={0}")
public static Iterable<Object[]> data() {
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkSchemaUpdateWithoutInputSchemaIT.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkSchemaUpdateWithoutInputSchemaIT.java
index d71d39a9c58..82963039790 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkSchemaUpdateWithoutInputSchemaIT.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkSchemaUpdateWithoutInputSchemaIT.java
@@ -17,6 +17,8 @@
*/
package org.apache.beam.sdk.io.gcp.bigquery;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
+
import java.io.IOException;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
import org.junit.AfterClass;
@@ -31,7 +33,7 @@ import org.junit.runners.Parameterized;
public class StorageApiSinkSchemaUpdateWithoutInputSchemaIT
extends StorageApiSinkSchemaUpdateITBase {
private static final String BIG_QUERY_DATASET_ID =
- "storage_api_sink_schema_change_without_input_" + System.nanoTime();
+ TEMP_DATASET_PREFIX + "sink_schema_change_without_input_" +
System.nanoTime();
@Parameterized.Parameters(name = "changeTableSchema={0}")
public static Iterable<Object[]> data() {
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/providers/BigQueryManagedIT.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/providers/BigQueryManagedIT.java
index 3106063b2d9..44fd7ecfe7f 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/providers/BigQueryManagedIT.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/providers/BigQueryManagedIT.java
@@ -17,6 +17,7 @@
*/
package org.apache.beam.sdk.io.gcp.bigquery.providers;
+import static
org.apache.beam.sdk.io.gcp.testing.BigtableUtils.TEMP_DATASET_PREFIX;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.hamcrest.Matchers.containsInAnyOrder;
@@ -83,7 +84,8 @@ public class BigQueryManagedIT {
private static final String PROJECT =
TestPipeline.testingPipelineOptions().as(GcpOptions.class).getProject();
- private static final String BIG_QUERY_DATASET_ID = "bigquery_managed_" +
System.nanoTime();
+ private static final String BIG_QUERY_DATASET_ID =
+ TEMP_DATASET_PREFIX + "managed_" + System.nanoTime();
private static final Clustering CLUSTERING = new
Clustering().setFields(Arrays.asList("str"));
diff --git a/sdks/python/apache_beam/io/gcp/bigquery_change_history_it_test.py
b/sdks/python/apache_beam/io/gcp/bigquery_change_history_it_test.py
index ef41fc393af..8cc7c7bb9c3 100644
--- a/sdks/python/apache_beam/io/gcp/bigquery_change_history_it_test.py
+++ b/sdks/python/apache_beam/io/gcp/bigquery_change_history_it_test.py
@@ -61,8 +61,8 @@ class BigQueryChangeHistoryIntegrationBase(unittest.TestCase):
cls.args = cls.test_pipeline.get_full_options_as_args()
cls.bq_wrapper = BigQueryWrapper()
suffix = secrets.token_hex(4)
- cls.dataset = f'beam_ch_src_{suffix}'
- cls.temp_dataset = f'beam_ch_tmp_{suffix}'
+ cls.dataset = f'{bigquery_tools._TEMP_DATASET_PREFIX}ch_src_{suffix}'
+ cls.temp_dataset = f'{bigquery_tools._TEMP_DATASET_PREFIX}ch_tmp_{suffix}'
cls.bq_wrapper.get_or_create_dataset(cls.project, cls.dataset)
ds = cls.bq_wrapper.client.datasets.Get(
bigquery.BigqueryDatasetsGetRequest(
diff --git a/sdks/python/apache_beam/io/gcp/bigquery_file_loads_test.py
b/sdks/python/apache_beam/io/gcp/bigquery_file_loads_test.py
index fda4a3e9d52..7ae0ccc8e32 100644
--- a/sdks/python/apache_beam/io/gcp/bigquery_file_loads_test.py
+++ b/sdks/python/apache_beam/io/gcp/bigquery_file_loads_test.py
@@ -1261,7 +1261,6 @@ class TestBigQueryFileLoads(_TestCaseWithTempDirCleanUp):
class BigQueryFileLoadsIT(unittest.TestCase):
- BIG_QUERY_DATASET_ID = 'python_bq_file_loads_'
BIG_QUERY_SCHEMA = (
'{"fields": [{"name": "name","type": "STRING"},'
'{"name": "language","type": "STRING"}]}')
@@ -1282,7 +1281,9 @@ class BigQueryFileLoadsIT(unittest.TestCase):
self.project = self.test_pipeline.get_option('project')
self.dataset_id = '%s%d%s' % (
- self.BIG_QUERY_DATASET_ID, int(time.time()), secrets.token_hex(3))
+ bigquery_tools._TEMP_DATASET_PREFIX,
+ int(time.time()),
+ secrets.token_hex(3))
self.bigquery_client = bigquery_tools.BigQueryWrapper()
self.bigquery_client.get_or_create_dataset(self.project, self.dataset_id)
self.output_table = "%s.output_table" % (self.dataset_id)
diff --git a/sdks/python/apache_beam/io/gcp/bigquery_read_it_test.py
b/sdks/python/apache_beam/io/gcp/bigquery_read_it_test.py
index 013bcc5d842..733b51a8052 100644
--- a/sdks/python/apache_beam/io/gcp/bigquery_read_it_test.py
+++ b/sdks/python/apache_beam/io/gcp/bigquery_read_it_test.py
@@ -88,8 +88,6 @@ def datetime_to_utc(element):
class BigQueryReadIntegrationTests(unittest.TestCase):
- BIG_QUERY_DATASET_ID = 'python_read_table_'
-
@classmethod
def setUpClass(cls):
cls.test_pipeline = TestPipeline(is_integration_test=True)
@@ -99,7 +97,9 @@ class BigQueryReadIntegrationTests(unittest.TestCase):
cls.bigquery_client = BigQueryWrapper()
cls.dataset_id = '%s%d%s' % (
- cls.BIG_QUERY_DATASET_ID, int(time.time()), secrets.token_hex(3))
+ bigquery_tools._TEMP_DATASET_PREFIX,
+ int(time.time()),
+ secrets.token_hex(3))
cls.bigquery_client.get_or_create_dataset(cls.project, cls.dataset_id)
_LOGGER.info(
"Created dataset %s in project %s", cls.dataset_id, cls.project)
@@ -309,7 +309,6 @@ class ReadTests(BigQueryReadIntegrationTests):
class ReadUsingStorageApiTests(BigQueryReadIntegrationTests):
- BIG_QUERY_DATASET_ID = 'python_read_table_'
TABLE_DATA = [{
'number': 1,
'string': '你好',
@@ -334,7 +333,9 @@ class
ReadUsingStorageApiTests(BigQueryReadIntegrationTests):
def setUpClass(cls):
super(ReadUsingStorageApiTests, cls).setUpClass()
cls.table_name = '%s%d%s' % (
- cls.BIG_QUERY_DATASET_ID, int(time.time()), secrets.token_hex(3))
+ bigquery_tools._TEMP_DATASET_PREFIX,
+ int(time.time()),
+ secrets.token_hex(3))
cls._create_table(cls.table_name)
table_id = '{}.{}'.format(cls.dataset_id, cls.table_name)
diff --git a/sdks/python/apache_beam/io/gcp/bigquery_test.py
b/sdks/python/apache_beam/io/gcp/bigquery_test.py
index dcadee7f6a1..874ae68362e 100644
--- a/sdks/python/apache_beam/io/gcp/bigquery_test.py
+++ b/sdks/python/apache_beam/io/gcp/bigquery_test.py
@@ -2953,15 +2953,15 @@ class PubSubBigQueryIT(unittest.TestCase):
@unittest.skipIf(HttpError is None, 'GCP dependencies are not installed')
class BigQueryFileLoadsIntegrationTests(unittest.TestCase):
- BIG_QUERY_DATASET_ID = 'python_bq_file_loads_'
-
def setUp(self):
self.test_pipeline = TestPipeline(is_integration_test=True)
self.runner_name = type(self.test_pipeline.runner).__name__
self.project = self.test_pipeline.get_option('project')
self.dataset_id = '%s%d%s' % (
- self.BIG_QUERY_DATASET_ID, int(time.time()), secrets.token_hex(3))
+ bigquery_tools._TEMP_DATASET_PREFIX,
+ int(time.time()),
+ secrets.token_hex(3))
self.bigquery_client = bigquery_tools.BigQueryWrapper()
self.bigquery_client.get_or_create_dataset(self.project, self.dataset_id)
self.output_table = '%s.output_table' % (self.dataset_id)
diff --git a/sdks/python/apache_beam/io/gcp/bigquery_tools.py
b/sdks/python/apache_beam/io/gcp/bigquery_tools.py
index 0d62ec5233c..5cc3441e171 100644
--- a/sdks/python/apache_beam/io/gcp/bigquery_tools.py
+++ b/sdks/python/apache_beam/io/gcp/bigquery_tools.py
@@ -109,6 +109,9 @@ _PROJECT_PATTERN = r'([a-z0-9.-]+:)?[a-z][a-z0-9-]*[a-z0-9]'
_DATASET_PATTERN = r'\w{1,1024}'
_TABLE_PATTERN = r'[\p{L}\p{M}\p{N}\p{Pc}\p{Pd}\p{Zs}$]{1,1024}'
+# For CI: temp dataset of this name are automatically deleted
+_TEMP_DATASET_PREFIX = 'beam_temp_dataset_'
+
# TODO(https://github.com/apache/beam/issues/25946): Add support for
# more Beam portable schema types as Python types
BIGQUERY_TYPE_TO_PYTHON_TYPE = {
diff --git
a/sdks/python/apache_beam/ml/inference/vertex_ai_model_monitoring_v2_it_test.py
b/sdks/python/apache_beam/ml/inference/vertex_ai_model_monitoring_v2_it_test.py
index f0903949752..e2a43ec6a3f 100644
---
a/sdks/python/apache_beam/ml/inference/vertex_ai_model_monitoring_v2_it_test.py
+++
b/sdks/python/apache_beam/ml/inference/vertex_ai_model_monitoring_v2_it_test.py
@@ -26,6 +26,7 @@ import uuid
import pytest
import apache_beam as beam
+from apache_beam.io.gcp import bigquery_tools
from apache_beam.ml.inference.base import ModelHandler
from apache_beam.ml.inference.base import PredictionResult
from apache_beam.ml.inference.base import RunInference
@@ -70,7 +71,7 @@ class
VertexAIModelMonitoringV2IntegrationTest(unittest.TestCase):
def test_vertex_ai_model_monitoring_v2_batch_pipeline(self):
test_pipeline = TestPipeline(is_integration_test=True)
job_id = str(uuid.uuid4())[:8]
- dataset_name = f"beam_mm_v2_{job_id}"
+ dataset_name = f"{bigquery_tools._TEMP_DATASET_PREFIX}mm_v2_{job_id}"
predictions_table_name = "predictions"
predictions_table_id =
f"{_ENDPOINT_PROJECT}:{dataset_name}.{predictions_table_name}"
display_name = f"beam-mm-v2-test-{job_id}"
@@ -278,7 +279,7 @@ class
VertexAIModelMonitoringV2IntegrationTest(unittest.TestCase):
test_pipeline = TestPipeline(
is_integration_test=True, additional_pipeline_args=["--streaming"])
job_id = str(uuid.uuid4())[:8]
- dataset_name = f"beam_mm_v2_str_{job_id}"
+ dataset_name = f"{bigquery_tools._TEMP_DATASET_PREFIX}mm_v2_str_{job_id}"
predictions_table_name = "predictions"
predictions_table_id =
f"{_ENDPOINT_PROJECT}:{dataset_name}.{predictions_table_name}"
display_name = f"beam-mm-v2-str-{job_id}"