diff --git a/it/google-cloud-platform/src/main/java/org/apache/beam/it/gcp/spanner/SpannerResourceManager.java b/it/google-cloud-platform/src/main/java/org/apache/beam/it/gcp/spanner/SpannerResourceManager.java index 62f5ebc7cc..f0343f24a5 100644 --- a/it/google-cloud-platform/src/main/java/org/apache/beam/it/gcp/spanner/SpannerResourceManager.java +++ b/it/google-cloud-platform/src/main/java/org/apache/beam/it/gcp/spanner/SpannerResourceManager.java @@ -63,6 +63,7 @@ import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.Optional; import java.util.Random; import java.util.function.Supplier; import java.util.stream.Collectors; @@ -794,10 +795,11 @@ public Builder useStaticInstance() { /** * Looks at the system properties if there's an instance id, and reuses it if configured. * + * @param staticInstanceHint optional instance Id of the spanner instance. * @return this builder with the instance ID set. */ @SuppressWarnings("nullness") - public Builder maybeUseStaticInstance() { + public Builder maybeUseStaticInstance(Optional staticInstanceHint) { String spannerInstanceId = System.getProperty("spannerInstanceId"); boolean isTestProject = Objects.equals(projectId, "cloud-teleport-testing") @@ -808,7 +810,12 @@ public Builder maybeUseStaticInstance() { if (isTestProject && shouldPickRandomInstance) { this.useStaticInstance = true; List staticInstanceList = TestConstants.SPANNER_TEST_INSTANCES; - this.instanceId = staticInstanceList.get(new Random().nextInt(staticInstanceList.size())); + if (staticInstanceHint.isPresent()) { + this.instanceId = + staticInstanceList.get(staticInstanceHint.get() % staticInstanceList.size()); + } else { + this.instanceId = staticInstanceList.get(new Random().nextInt(staticInstanceList.size())); + } } else if (spannerInstanceId != null) { this.useStaticInstance = true; this.instanceId = spannerInstanceId; @@ -817,6 +824,17 @@ public Builder maybeUseStaticInstance() { return this; } + /** + * Looks at the system properties if there's an instance id, and reuses it if configured. + * + * @return this builder with the instance ID set. + */ + @SuppressWarnings("nullness") + public Builder maybeUseStaticInstance() { + maybeUseStaticInstance(Optional.empty()); + return this; + } + /** * Set the instance ID of a static Spanner instance for this Resource Manager to manage. * diff --git a/v2/datastream-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/DataStreamToSpannerLTBase.java b/v2/datastream-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/DataStreamToSpannerLTBase.java index f22e44dfc0..75713d6849 100644 --- a/v2/datastream-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/DataStreamToSpannerLTBase.java +++ b/v2/datastream-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/DataStreamToSpannerLTBase.java @@ -40,6 +40,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.concurrent.Callable; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -113,7 +114,7 @@ public void setUpResourceManagers(String spannerDdlResource, boolean separateSha testRootDir = getClass().getSimpleName(); spannerResourceManager = SpannerResourceManager.builder(testName, project, region) - .maybeUseStaticInstance() + .maybeUseStaticInstance(Optional.of(1)) .setNodeCount(10) .setMonitoringClient(monitoringClient) .setSuppressVerboseLogs(true) @@ -138,7 +139,7 @@ public void setUpResourceManagers(String spannerDdlResource, boolean separateSha if (separateShadowTableDb) { shadowTableSpannerResourceManager = SpannerResourceManager.builder("shadow_" + testName, project, region) - .maybeUseStaticInstance() + .maybeUseStaticInstance(Optional.of(1)) .setNodeCount(10) .setMonitoringClient(monitoringClient) .setSuppressVerboseLogs(true) diff --git a/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/MySQL5KTablesLT.java b/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/MySQL5KTablesLT.java index 06d1d43f18..10a604f6cf 100644 --- a/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/MySQL5KTablesLT.java +++ b/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/MySQL5KTablesLT.java @@ -34,6 +34,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; @@ -84,7 +85,7 @@ public void setUp() throws IOException { mySQLResourceManager = MySQLResourceManager.builder(testName).build(); spannerResourceManager = SpannerResourceManager.builder(testName, project, region) - .maybeUseStaticInstance() + .maybeUseStaticInstance(Optional.of(3)) .setMonitoringClient(monitoringClient) .build(); diff --git a/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/MySQLMultiSharded1024ShardsLT.java b/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/MySQLMultiSharded1024ShardsLT.java index 5442d3c030..2c5c93f1ea 100644 --- a/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/MySQLMultiSharded1024ShardsLT.java +++ b/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/MySQLMultiSharded1024ShardsLT.java @@ -36,6 +36,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; @@ -94,7 +95,7 @@ public void setUp() throws IOException { spannerResourceManager = SpannerResourceManager.builder(testName, project, region) - .maybeUseStaticInstance() + .maybeUseStaticInstance(Optional.of(3)) .setMonitoringClient(monitoringClient) .build(); diff --git a/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/PostgreSQL5KTablesLT.java b/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/PostgreSQL5KTablesLT.java index 8786db1662..515d410656 100644 --- a/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/PostgreSQL5KTablesLT.java +++ b/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/PostgreSQL5KTablesLT.java @@ -34,6 +34,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; @@ -79,7 +80,7 @@ public void setUp() throws IOException { postgresResourceManager = PostgresResourceManager.builder(testName).build(); spannerResourceManager = SpannerResourceManager.builder(testName, project, region) - .maybeUseStaticInstance() + .maybeUseStaticInstance(Optional.of(3)) .setMonitoringClient(monitoringClient) .build(); diff --git a/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/PostgreSQLMultiSharded1024ShardsLT.java b/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/PostgreSQLMultiSharded1024ShardsLT.java index ed760ba441..c81006059e 100644 --- a/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/PostgreSQLMultiSharded1024ShardsLT.java +++ b/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/PostgreSQLMultiSharded1024ShardsLT.java @@ -36,6 +36,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; @@ -95,7 +96,7 @@ public void setUp() throws IOException { spannerResourceManager = SpannerResourceManager.builder(testName, project, region) - .maybeUseStaticInstance() + .maybeUseStaticInstance(Optional.of(3)) .setMonitoringClient(monitoringClient) .build(); diff --git a/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/SourceDbToSpannerLTBase.java b/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/SourceDbToSpannerLTBase.java index 3313da4cc5..d39cde4ffb 100644 --- a/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/SourceDbToSpannerLTBase.java +++ b/v2/sourcedb-to-spanner/src/test/java/com/google/cloud/teleport/v2/templates/loadtesting/SourceDbToSpannerLTBase.java @@ -31,6 +31,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.stream.Collectors; import org.apache.beam.it.common.PipelineLauncher; import org.apache.beam.it.common.PipelineLauncher.LaunchConfig; @@ -111,7 +112,7 @@ public void setUp( spannerResourceManager = SpannerResourceManager.builder(testName, project, region) - .maybeUseStaticInstance() + .maybeUseStaticInstance(Optional.of(3)) .setNodeCount(SPANNER_NODE_COUNT) .setMonitoringClient(monitoringClient) .build(); diff --git a/v2/spanner-to-sourcedb/src/test/java/com/google/cloud/teleport/v2/templates/SpannerToSourceDb5kTablesLT.java b/v2/spanner-to-sourcedb/src/test/java/com/google/cloud/teleport/v2/templates/SpannerToSourceDb5kTablesLT.java index 666c1ac730..eb99b5a5ba 100644 --- a/v2/spanner-to-sourcedb/src/test/java/com/google/cloud/teleport/v2/templates/SpannerToSourceDb5kTablesLT.java +++ b/v2/spanner-to-sourcedb/src/test/java/com/google/cloud/teleport/v2/templates/SpannerToSourceDb5kTablesLT.java @@ -35,6 +35,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.concurrent.Callable; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -84,21 +85,21 @@ public void setUp() throws IOException { spannerResourceManager = SpannerResourceManager.builder("rr-main-" + testName, project, region) - .maybeUseStaticInstance() + .maybeUseStaticInstance(Optional.of(2)) .setMonitoringClient(monitoringClient) .setSuppressVerboseLogs(true) .build(); spannerMetadataResourceManager = SpannerResourceManager.builder("rr-meta-" + testName, project, region) - .maybeUseStaticInstance() + .maybeUseStaticInstance(Optional.of(2)) .setSuppressVerboseLogs(true) .build(); spannerMetadataResourceManager.ensureUsableAndCreateResources(); spannerChangeStreamMetadataResourceManager = SpannerResourceManager.builder("rr-cs-meta-" + testName, project, region) - .maybeUseStaticInstance() + .maybeUseStaticInstance(Optional.of(2)) .setSuppressVerboseLogs(true) .build(); spannerChangeStreamMetadataResourceManager.ensureUsableAndCreateResources(); diff --git a/v2/spanner-to-sourcedb/src/test/java/com/google/cloud/teleport/v2/templates/SpannerToSourceDbLTBase.java b/v2/spanner-to-sourcedb/src/test/java/com/google/cloud/teleport/v2/templates/SpannerToSourceDbLTBase.java index 88a89bcbf1..cac9550cbf 100644 --- a/v2/spanner-to-sourcedb/src/test/java/com/google/cloud/teleport/v2/templates/SpannerToSourceDbLTBase.java +++ b/v2/spanner-to-sourcedb/src/test/java/com/google/cloud/teleport/v2/templates/SpannerToSourceDbLTBase.java @@ -33,6 +33,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.Random; import org.apache.beam.it.common.PipelineLauncher; import org.apache.beam.it.common.PipelineLauncher.LaunchConfig; @@ -136,7 +137,7 @@ public SpannerResourceManager createSpannerDatabase(String spannerDdlResourceFil throws IOException { SpannerResourceManager spannerResourceManager = SpannerResourceManager.builder("rr-loadtest-" + testName, project, region) - .maybeUseStaticInstance() + .maybeUseStaticInstance(Optional.of(2)) .build(); String ddl = String.join( @@ -161,7 +162,7 @@ public SpannerResourceManager createSpannerMetadataDatabase() throws IOException if (metadataInstanceId != null && !metadataInstanceId.isEmpty()) { builder.setInstanceId(metadataInstanceId).useStaticInstance(); } else { - builder.maybeUseStaticInstance(); + builder.maybeUseStaticInstance(Optional.of(2)); } SpannerResourceManager spannerMetadataResourceManager = builder.build();