From 16b5244435e4118a6324060e253cb73f9a2a54a1 Mon Sep 17 00:00:00 2001 From: Steinar Eliassen Date: Tue, 2 Sep 2025 08:58:00 +0200 Subject: [PATCH 1/4] Expose gRPC endpoint --- .../testcontainers/containers/BigQueryEmulatorContainer.java | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/modules/gcloud/src/main/java/org/testcontainers/containers/BigQueryEmulatorContainer.java b/modules/gcloud/src/main/java/org/testcontainers/containers/BigQueryEmulatorContainer.java index 6c305a50537..cc731ba7426 100644 --- a/modules/gcloud/src/main/java/org/testcontainers/containers/BigQueryEmulatorContainer.java +++ b/modules/gcloud/src/main/java/org/testcontainers/containers/BigQueryEmulatorContainer.java @@ -33,6 +33,10 @@ public String getEmulatorHttpEndpoint() { return String.format("http://%s:%d", getHost(), getMappedPort(HTTP_PORT)); } + public String getEmulatorGrpcEndpoint() { + return String.format("http://%s:%d", getHost(), getMappedPort(GRPC_PORT)); + } + public String getProjectId() { return PROJECT_ID; } From b9dad59db20f820eea28b7f289f29eb78417a497 Mon Sep 17 00:00:00 2001 From: Steinar Eliassen Date: Tue, 9 Sep 2025 19:09:02 +0200 Subject: [PATCH 2/4] Modify exposing endpoint to exposing port, only port is required --- .../containers/BigQueryEmulatorContainer.java | 4 +- .../BigQueryEmulatorContainerTest.java | 56 +++++++++++++++++++ 2 files changed, 58 insertions(+), 2 deletions(-) diff --git a/modules/gcloud/src/main/java/org/testcontainers/containers/BigQueryEmulatorContainer.java b/modules/gcloud/src/main/java/org/testcontainers/containers/BigQueryEmulatorContainer.java index cc731ba7426..6831785b0a5 100644 --- a/modules/gcloud/src/main/java/org/testcontainers/containers/BigQueryEmulatorContainer.java +++ b/modules/gcloud/src/main/java/org/testcontainers/containers/BigQueryEmulatorContainer.java @@ -33,8 +33,8 @@ public String getEmulatorHttpEndpoint() { return String.format("http://%s:%d", getHost(), getMappedPort(HTTP_PORT)); } - public String getEmulatorGrpcEndpoint() { - return String.format("http://%s:%d", getHost(), getMappedPort(GRPC_PORT)); + public Integer getEmulatorGrpcPort() { + return getMappedPort(GRPC_PORT); } public String getProjectId() { diff --git a/modules/gcloud/src/test/java/org/testcontainers/containers/BigQueryEmulatorContainerTest.java b/modules/gcloud/src/test/java/org/testcontainers/containers/BigQueryEmulatorContainerTest.java index 7f06a67c26c..4374c8a1378 100644 --- a/modules/gcloud/src/test/java/org/testcontainers/containers/BigQueryEmulatorContainerTest.java +++ b/modules/gcloud/src/test/java/org/testcontainers/containers/BigQueryEmulatorContainerTest.java @@ -1,12 +1,23 @@ package org.testcontainers.containers; +import com.google.api.gax.core.NoCredentialsProvider; +import com.google.api.gax.grpc.GrpcTransportChannel; +import com.google.api.gax.rpc.FixedTransportChannelProvider; import com.google.cloud.NoCredentials; import com.google.cloud.bigquery.BigQuery; import com.google.cloud.bigquery.BigQueryOptions; import com.google.cloud.bigquery.QueryJobConfiguration; import com.google.cloud.bigquery.TableResult; +import com.google.cloud.bigquery.storage.v1.BigQueryWriteClient; +import com.google.cloud.bigquery.storage.v1.BigQueryWriteSettings; +import com.google.cloud.bigquery.storage.v1.CreateWriteStreamRequest; +import com.google.cloud.bigquery.storage.v1.TableName; +import com.google.cloud.bigquery.storage.v1.WriteStream; +import io.grpc.ManagedChannelBuilder; +import org.threeten.bp.Duration; import org.junit.jupiter.api.Test; +import java.io.IOException; import java.math.BigDecimal; import java.util.List; import java.util.stream.Collectors; @@ -15,6 +26,51 @@ class BigQueryEmulatorContainerTest { + @Test + public void testGrcp() throws IOException { + // Shallow test, validate that connection can be set up, and attempt to create write stream fails. + // BigQueryWriteSettings requires a HTTP/2 connection, not provided by the originally exposed endpoint. A "not found" exceptionm + // indicates successful + try (BigQueryEmulatorContainer container = new BigQueryEmulatorContainer("ghcr.io/goccy/bigquery-emulator:0.6.5")) { + container.start(); + BigQueryWriteSettings.Builder bigQueryWriteSettingsBuilder = BigQueryWriteSettings.newBuilder(); + + bigQueryWriteSettingsBuilder.createWriteStreamSettings() + .setRetrySettings(bigQueryWriteSettingsBuilder.createWriteStreamSettings() + .getRetrySettings() + .toBuilder() + .setTotalTimeout(Duration.ofSeconds(60)) + .build()); + + BigQueryWriteClient bigQueryWriteClient = BigQueryWriteClient.create( + bigQueryWriteSettingsBuilder.setTransportChannelProvider(FixedTransportChannelProvider.create(GrpcTransportChannel.create( + ManagedChannelBuilder.forAddress(container.getHost(), container.getEmulatorGrpcPort()).usePlaintext().build()))) + .setCredentialsProvider(NoCredentialsProvider.create()) + .build() + ); + + TableName parentTable = TableName.of(container.getProjectId(), "dataset", "table"); + CreateWriteStreamRequest createWriteStreamRequest = CreateWriteStreamRequest.newBuilder() + .setParent(parentTable.toString()) + .setWriteStream(WriteStream.newBuilder().setType(WriteStream.Type.PENDING)) + .build(); + + String message = null; + try { + // This will fail, extract error message to check that it fails in a "we reached the backend" way to ensure that setup was correct + WriteStream writeStream = bigQueryWriteClient.createWriteStream(createWriteStreamRequest); + // Example setting up StreamWriter. Note passing bigQueryWriteClient as parameter, this is needed to avoid using gcloud credentials: + /* StreamWriter writer = StreamWriter.newBuilder(writeStream.getName(), bigQueryWriteClient).setWriterSchema(schema).build(); */ + } catch (RuntimeException e) { + message = e.getMessage(); + } + assertThat(message).contains("dataset dataset is not found in project test-project"); + + bigQueryWriteClient.shutdown(); + bigQueryWriteClient.close(); + } + } + @Test void test() throws Exception { try ( From ca5fbd7e922749c8ce34e3b27aea069c8a325096 Mon Sep 17 00:00:00 2001 From: Steinar Eliassen Date: Mon, 15 Sep 2025 16:36:22 +0200 Subject: [PATCH 3/4] Expand test to include a proper table we can connect to, and validate that we can connect a BigQueryWriteClient to the table. --- .../BigQueryEmulatorContainerTest.java | 130 +++++++++++------- 1 file changed, 77 insertions(+), 53 deletions(-) diff --git a/modules/gcloud/src/test/java/org/testcontainers/containers/BigQueryEmulatorContainerTest.java b/modules/gcloud/src/test/java/org/testcontainers/containers/BigQueryEmulatorContainerTest.java index 4374c8a1378..45d8d617031 100644 --- a/modules/gcloud/src/test/java/org/testcontainers/containers/BigQueryEmulatorContainerTest.java +++ b/modules/gcloud/src/test/java/org/testcontainers/containers/BigQueryEmulatorContainerTest.java @@ -6,7 +6,16 @@ import com.google.cloud.NoCredentials; import com.google.cloud.bigquery.BigQuery; import com.google.cloud.bigquery.BigQueryOptions; +import com.google.cloud.bigquery.DatasetId; +import com.google.cloud.bigquery.DatasetInfo; +import com.google.cloud.bigquery.Field; import com.google.cloud.bigquery.QueryJobConfiguration; +import com.google.cloud.bigquery.Schema; +import com.google.cloud.bigquery.StandardSQLTypeName; +import com.google.cloud.bigquery.StandardTableDefinition; +import com.google.cloud.bigquery.TableDefinition; +import com.google.cloud.bigquery.TableId; +import com.google.cloud.bigquery.TableInfo; import com.google.cloud.bigquery.TableResult; import com.google.cloud.bigquery.storage.v1.BigQueryWriteClient; import com.google.cloud.bigquery.storage.v1.BigQueryWriteSettings; @@ -26,53 +35,19 @@ class BigQueryEmulatorContainerTest { - @Test - public void testGrcp() throws IOException { - // Shallow test, validate that connection can be set up, and attempt to create write stream fails. - // BigQueryWriteSettings requires a HTTP/2 connection, not provided by the originally exposed endpoint. A "not found" exceptionm - // indicates successful - try (BigQueryEmulatorContainer container = new BigQueryEmulatorContainer("ghcr.io/goccy/bigquery-emulator:0.6.5")) { - container.start(); - BigQueryWriteSettings.Builder bigQueryWriteSettingsBuilder = BigQueryWriteSettings.newBuilder(); - - bigQueryWriteSettingsBuilder.createWriteStreamSettings() - .setRetrySettings(bigQueryWriteSettingsBuilder.createWriteStreamSettings() - .getRetrySettings() - .toBuilder() - .setTotalTimeout(Duration.ofSeconds(60)) - .build()); - - BigQueryWriteClient bigQueryWriteClient = BigQueryWriteClient.create( - bigQueryWriteSettingsBuilder.setTransportChannelProvider(FixedTransportChannelProvider.create(GrpcTransportChannel.create( - ManagedChannelBuilder.forAddress(container.getHost(), container.getEmulatorGrpcPort()).usePlaintext().build()))) - .setCredentialsProvider(NoCredentialsProvider.create()) - .build() - ); - - TableName parentTable = TableName.of(container.getProjectId(), "dataset", "table"); - CreateWriteStreamRequest createWriteStreamRequest = CreateWriteStreamRequest.newBuilder() - .setParent(parentTable.toString()) - .setWriteStream(WriteStream.newBuilder().setType(WriteStream.Type.PENDING)) - .build(); - - String message = null; - try { - // This will fail, extract error message to check that it fails in a "we reached the backend" way to ensure that setup was correct - WriteStream writeStream = bigQueryWriteClient.createWriteStream(createWriteStreamRequest); - // Example setting up StreamWriter. Note passing bigQueryWriteClient as parameter, this is needed to avoid using gcloud credentials: - /* StreamWriter writer = StreamWriter.newBuilder(writeStream.getName(), bigQueryWriteClient).setWriterSchema(schema).build(); */ - } catch (RuntimeException e) { - message = e.getMessage(); - } - assertThat(message).contains("dataset dataset is not found in project test-project"); - - bigQueryWriteClient.shutdown(); - bigQueryWriteClient.close(); - } + private BigQuery getBigQuery(BigQueryEmulatorContainer container) { + String url = container.getEmulatorHttpEndpoint(); + return BigQueryOptions + .newBuilder() + .setProjectId(container.getProjectId()) + .setHost(url) + .setLocation(url) + .setCredentials(NoCredentials.getInstance()) + .build().getService(); } @Test - void test() throws Exception { + void testHttpEndpoint() throws Exception { try ( // emulatorContainer { BigQueryEmulatorContainer container = new BigQueryEmulatorContainer("ghcr.io/goccy/bigquery-emulator:0.4.3") @@ -81,15 +56,7 @@ void test() throws Exception { container.start(); // bigQueryClient { - String url = container.getEmulatorHttpEndpoint(); - BigQueryOptions options = BigQueryOptions - .newBuilder() - .setProjectId(container.getProjectId()) - .setHost(url) - .setLocation(url) - .setCredentials(NoCredentials.getInstance()) - .build(); - BigQuery bigQuery = options.getService(); + BigQuery bigQuery = getBigQuery(container); // } String fn = @@ -107,4 +74,61 @@ void test() throws Exception { assertThat(values).containsOnly(BigDecimal.valueOf(30)); } } + + @Test + void testGrcpEndpoint() throws IOException { + try (BigQueryEmulatorContainer container = new BigQueryEmulatorContainer("ghcr.io/goccy/bigquery-emulator:0.6.5")) { + container.start(); + + // Test setup. + // Create a table the "regular" way. We need this to verify we can connect a writestream + BigQuery bigQuery = getBigQuery(container); + String tableName = "test-table"; + String datasetName = "test-dataset"; + + bigQuery.create(DatasetInfo.of(DatasetId.of(container.getProjectId(), datasetName))); + + Schema schema = Schema.of( + Field.of("name", StandardSQLTypeName.STRING) + ); + + TableId tableId = TableId.of(datasetName, tableName); + TableDefinition tableDefinition = StandardTableDefinition.of(schema); + TableInfo tableInfo = TableInfo.newBuilder(tableId, tableDefinition).build(); + + bigQuery.create(tableInfo); + + // Actual test. + // BigQueryWriteSettings requires a HTTP/2 connection, not provided by the originally exposed endpoint. + BigQueryWriteSettings.Builder bigQueryWriteSettingsBuilder = BigQueryWriteSettings.newBuilder(); + + bigQueryWriteSettingsBuilder.createWriteStreamSettings() + .setRetrySettings(bigQueryWriteSettingsBuilder.createWriteStreamSettings() + .getRetrySettings() + .toBuilder() + .setTotalTimeout(Duration.ofSeconds(60)) + .build()); + + // Use the now exposed grpcPort to get a working connection. + BigQueryWriteClient bigQueryWriteClient = BigQueryWriteClient.create( + bigQueryWriteSettingsBuilder.setTransportChannelProvider(FixedTransportChannelProvider.create(GrpcTransportChannel.create( + ManagedChannelBuilder.forAddress(container.getHost(), container.getEmulatorGrpcPort()).usePlaintext().build()))) + .setCredentialsProvider(NoCredentialsProvider.create()) + .build() + ); + + TableName parentTable = TableName.of(container.getProjectId(), datasetName, tableName); + CreateWriteStreamRequest createWriteStreamRequest = CreateWriteStreamRequest.newBuilder() + .setParent(parentTable.toString()) + .setWriteStream(WriteStream.newBuilder().setType(WriteStream.Type.PENDING)) + .build(); + + // Validate that we can successfully create a write stream. This would not work with http endpoint + bigQueryWriteClient.createWriteStream(createWriteStreamRequest); + + bigQueryWriteClient.shutdown(); + bigQueryWriteClient.close(); + } + } + } From ba27f37fdb5397731c9ca45f01657005b8cdb0b4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Edd=C3=BA=20Mel=C3=A9ndez?= Date: Wed, 24 Sep 2025 12:07:07 -0600 Subject: [PATCH 4/4] Polish --- .../BigQueryEmulatorContainerTest.java | 142 +++++++++++++----- 1 file changed, 106 insertions(+), 36 deletions(-) diff --git a/modules/gcloud/src/test/java/org/testcontainers/containers/BigQueryEmulatorContainerTest.java b/modules/gcloud/src/test/java/org/testcontainers/containers/BigQueryEmulatorContainerTest.java index 45d8d617031..0337cd563b9 100644 --- a/modules/gcloud/src/test/java/org/testcontainers/containers/BigQueryEmulatorContainerTest.java +++ b/modules/gcloud/src/test/java/org/testcontainers/containers/BigQueryEmulatorContainerTest.java @@ -1,5 +1,6 @@ package org.testcontainers.containers; +import com.google.api.core.ApiFuture; import com.google.api.gax.core.NoCredentialsProvider; import com.google.api.gax.grpc.GrpcTransportChannel; import com.google.api.gax.rpc.FixedTransportChannelProvider; @@ -17,16 +18,23 @@ import com.google.cloud.bigquery.TableId; import com.google.cloud.bigquery.TableInfo; import com.google.cloud.bigquery.TableResult; +import com.google.cloud.bigquery.storage.v1.AppendRowsResponse; +import com.google.cloud.bigquery.storage.v1.BatchCommitWriteStreamsRequest; +import com.google.cloud.bigquery.storage.v1.BatchCommitWriteStreamsResponse; import com.google.cloud.bigquery.storage.v1.BigQueryWriteClient; import com.google.cloud.bigquery.storage.v1.BigQueryWriteSettings; import com.google.cloud.bigquery.storage.v1.CreateWriteStreamRequest; +import com.google.cloud.bigquery.storage.v1.FinalizeWriteStreamRequest; +import com.google.cloud.bigquery.storage.v1.FinalizeWriteStreamResponse; +import com.google.cloud.bigquery.storage.v1.JsonStreamWriter; import com.google.cloud.bigquery.storage.v1.TableName; import com.google.cloud.bigquery.storage.v1.WriteStream; import io.grpc.ManagedChannelBuilder; -import org.threeten.bp.Duration; +import org.json.JSONArray; +import org.json.JSONObject; import org.junit.jupiter.api.Test; +import org.threeten.bp.Duration; -import java.io.IOException; import java.math.BigDecimal; import java.util.List; import java.util.stream.Collectors; @@ -35,17 +43,6 @@ class BigQueryEmulatorContainerTest { - private BigQuery getBigQuery(BigQueryEmulatorContainer container) { - String url = container.getEmulatorHttpEndpoint(); - return BigQueryOptions - .newBuilder() - .setProjectId(container.getProjectId()) - .setHost(url) - .setLocation(url) - .setCredentials(NoCredentials.getInstance()) - .build().getService(); - } - @Test void testHttpEndpoint() throws Exception { try ( @@ -56,7 +53,15 @@ void testHttpEndpoint() throws Exception { container.start(); // bigQueryClient { - BigQuery bigQuery = getBigQuery(container); + String url = container.getEmulatorHttpEndpoint(); + BigQueryOptions options = BigQueryOptions + .newBuilder() + .setProjectId(container.getProjectId()) + .setHost(url) + .setLocation(url) + .setCredentials(NoCredentials.getInstance()) + .build(); + BigQuery bigQuery = options.getService(); // } String fn = @@ -76,21 +81,19 @@ void testHttpEndpoint() throws Exception { } @Test - void testGrcpEndpoint() throws IOException { - try (BigQueryEmulatorContainer container = new BigQueryEmulatorContainer("ghcr.io/goccy/bigquery-emulator:0.6.5")) { + void testGrcpEndpoint() throws Exception { + try ( + BigQueryEmulatorContainer container = new BigQueryEmulatorContainer("ghcr.io/goccy/bigquery-emulator:0.6.5") + ) { container.start(); - // Test setup. - // Create a table the "regular" way. We need this to verify we can connect a writestream - BigQuery bigQuery = getBigQuery(container); + BigQuery bigQuery = getBigQuery(container); String tableName = "test-table"; String datasetName = "test-dataset"; bigQuery.create(DatasetInfo.of(DatasetId.of(container.getProjectId(), datasetName))); - Schema schema = Schema.of( - Field.of("name", StandardSQLTypeName.STRING) - ); + Schema schema = Schema.of(Field.of("name", StandardSQLTypeName.STRING)); TableId tableId = TableId.of(datasetName, tableName); TableDefinition tableDefinition = StandardTableDefinition.of(schema); @@ -98,37 +101,104 @@ void testGrcpEndpoint() throws IOException { bigQuery.create(tableInfo); - // Actual test. - // BigQueryWriteSettings requires a HTTP/2 connection, not provided by the originally exposed endpoint. BigQueryWriteSettings.Builder bigQueryWriteSettingsBuilder = BigQueryWriteSettings.newBuilder(); - bigQueryWriteSettingsBuilder.createWriteStreamSettings() - .setRetrySettings(bigQueryWriteSettingsBuilder.createWriteStreamSettings() - .getRetrySettings() - .toBuilder() - .setTotalTimeout(Duration.ofSeconds(60)) - .build()); + bigQueryWriteSettingsBuilder + .createWriteStreamSettings() + .setRetrySettings( + bigQueryWriteSettingsBuilder + .createWriteStreamSettings() + .getRetrySettings() + .toBuilder() + .setTotalTimeout(Duration.ofSeconds(60)) + .build() + ); - // Use the now exposed grpcPort to get a working connection. BigQueryWriteClient bigQueryWriteClient = BigQueryWriteClient.create( - bigQueryWriteSettingsBuilder.setTransportChannelProvider(FixedTransportChannelProvider.create(GrpcTransportChannel.create( - ManagedChannelBuilder.forAddress(container.getHost(), container.getEmulatorGrpcPort()).usePlaintext().build()))) + bigQueryWriteSettingsBuilder + .setTransportChannelProvider( + FixedTransportChannelProvider.create( + GrpcTransportChannel.create( + ManagedChannelBuilder + .forAddress(container.getHost(), container.getEmulatorGrpcPort()) + .usePlaintext() + .build() + ) + ) + ) .setCredentialsProvider(NoCredentialsProvider.create()) .build() ); TableName parentTable = TableName.of(container.getProjectId(), datasetName, tableName); - CreateWriteStreamRequest createWriteStreamRequest = CreateWriteStreamRequest.newBuilder() + CreateWriteStreamRequest createWriteStreamRequest = CreateWriteStreamRequest + .newBuilder() .setParent(parentTable.toString()) .setWriteStream(WriteStream.newBuilder().setType(WriteStream.Type.PENDING)) .build(); - // Validate that we can successfully create a write stream. This would not work with http endpoint - bigQueryWriteClient.createWriteStream(createWriteStreamRequest); + WriteStream writeStream = bigQueryWriteClient.createWriteStream(createWriteStreamRequest); + + JsonStreamWriter writer = JsonStreamWriter + .newBuilder(writeStream.getName(), writeStream.getTableSchema(), bigQueryWriteClient) + .build(); + + JSONArray jsonArray = new JSONArray(); + JSONObject record1 = new JSONObject(); + record1.put("name", "Alice"); + jsonArray.put(record1); + + JSONObject record2 = new JSONObject(); + record2.put("name", "Bob"); + jsonArray.put(record2); + + ApiFuture future = writer.append(jsonArray); + AppendRowsResponse response = future.get(); + + FinalizeWriteStreamRequest finalizeRequest = FinalizeWriteStreamRequest + .newBuilder() + .setName(writeStream.getName()) + .build(); + FinalizeWriteStreamResponse finalizeResponse = bigQueryWriteClient.finalizeWriteStream(finalizeRequest); + + BatchCommitWriteStreamsRequest commitRequest = BatchCommitWriteStreamsRequest + .newBuilder() + .setParent(parentTable.toString()) + .addWriteStreams(writeStream.getName()) + .build(); + BatchCommitWriteStreamsResponse commitResponse = bigQueryWriteClient.batchCommitWriteStreams(commitRequest); + + writer.close(); + + String sql = String.format( + "SELECT name FROM `%s.%s.%s` ORDER BY name", + container.getProjectId(), + datasetName, + tableName + ); + TableResult result = bigQuery.query(QueryJobConfiguration.newBuilder(sql).build()); + + List names = result + .streamValues() + .map(row -> row.get("name").getStringValue()) + .collect(Collectors.toList()); + + assertThat(names).containsExactly("Alice", "Bob"); bigQueryWriteClient.shutdown(); bigQueryWriteClient.close(); } } + private BigQuery getBigQuery(BigQueryEmulatorContainer container) { + String url = container.getEmulatorHttpEndpoint(); + return BigQueryOptions + .newBuilder() + .setProjectId(container.getProjectId()) + .setHost(url) + .setLocation(url) + .setCredentials(NoCredentials.getInstance()) + .build() + .getService(); + } }