This is an automated email from the ASF dual-hosted git repository.
yhu 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 fa43f82f8a2 Disable BigQueryStorageStreamSource.splitAtFraction when
read api v2 used (#30443)
fa43f82f8a2 is described below
commit fa43f82f8a2714ce1621dd9f70af734f1196be54
Author: Yi Hu <[email protected]>
AuthorDate: Mon Mar 4 18:40:24 2024 -0500
Disable BigQueryStorageStreamSource.splitAtFraction when read api v2 used
(#30443)
---
.../gcp/bigquery/BigQueryStorageStreamSource.java | 5 +++-
.../io/gcp/bigquery/BigQueryIOStorageReadTest.java | 35 ++++++++++++++++++++++
2 files changed, 39 insertions(+), 1 deletion(-)
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryStorageStreamSource.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryStorageStreamSource.java
index 8f7f50febaf..436a00a6b77 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryStorageStreamSource.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryStorageStreamSource.java
@@ -165,6 +165,7 @@ class BigQueryStorageStreamSource<T> extends
BoundedSource<T> {
// Values used for progress reporting.
private boolean splitPossible = true;
+ private boolean splitAllowed = true;
private double fractionConsumed;
private double progressAtResponseStart;
private double progressAtResponseEnd;
@@ -199,6 +200,8 @@ class BigQueryStorageStreamSource<T> extends
BoundedSource<T> {
this.parseFn = source.parseFn;
this.storageClient = source.bqServices.getStorageClient(options);
this.tableSchema = fromJsonString(source.jsonTableSchema,
TableSchema.class);
+ // number of stream determined from server side for storage read api v2
+ this.splitAllowed = !options.getEnableStorageReadApiV2();
this.fractionConsumed = 0d;
this.progressAtResponseStart = 0d;
this.progressAtResponseEnd = 0d;
@@ -341,7 +344,7 @@ class BigQueryStorageStreamSource<T> extends
BoundedSource<T> {
return null;
}
- if (!splitPossible) {
+ if (!splitPossible || !splitAllowed) {
return null;
}
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOStorageReadTest.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOStorageReadTest.java
index 5a78f529e8f..a23dd866eea 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOStorageReadTest.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOStorageReadTest.java
@@ -95,11 +95,13 @@ import
org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO.TableRowParser;
import org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO.TypedRead;
import org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO.TypedRead.Method;
import org.apache.beam.sdk.io.gcp.bigquery.BigQueryServices.StorageClient;
+import
org.apache.beam.sdk.io.gcp.bigquery.BigQueryStorageStreamSource.BigQueryStorageStreamReader;
import org.apache.beam.sdk.io.gcp.bigquery.BigQueryUtils.ConversionOptions;
import org.apache.beam.sdk.io.gcp.testing.FakeBigQueryServices;
import
org.apache.beam.sdk.io.gcp.testing.FakeBigQueryServices.FakeBigQueryServerStream;
import org.apache.beam.sdk.io.gcp.testing.FakeDatasetService;
import org.apache.beam.sdk.options.PipelineOptions;
+import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.options.ValueProvider;
import org.apache.beam.sdk.options.ValueProvider.StaticValueProvider;
import org.apache.beam.sdk.schemas.FieldAccessDescriptor;
@@ -767,6 +769,39 @@ public class BigQueryIOStorageReadTest {
assertThat(streamSource.split(0, options),
containsInAnyOrder(streamSource));
}
+ @Test
+ public void testSplitReadStreamAtFraction() throws IOException {
+
+ ReadSession readSession =
+ ReadSession.newBuilder()
+ .setName("readSession")
+
.setAvroSchema(AvroSchema.newBuilder().setSchema(AVRO_SCHEMA_STRING))
+ .build();
+
+ ReadRowsRequest expectedRequest =
+ ReadRowsRequest.newBuilder().setReadStream("readStream").build();
+ List<ReadRowsResponse> responses = ImmutableList.of();
+
+ StorageClient fakeStorageClient = mock(StorageClient.class);
+ when(fakeStorageClient.readRows(expectedRequest, ""))
+ .thenReturn(new FakeBigQueryServerStream<>(responses));
+
+ BigQueryStorageStreamSource<TableRow> streamSource =
+ BigQueryStorageStreamSource.create(
+ readSession,
+ ReadStream.newBuilder().setName("readStream").build(),
+ TABLE_SCHEMA,
+ new TableRowParser(),
+ TableRowJsonCoder.of(),
+ new FakeBigQueryServices().withStorageClient(fakeStorageClient));
+
+ PipelineOptions options =
PipelineOptionsFactory.fromArgs("--enableStorageReadApiV2").create();
+ BigQueryStorageStreamReader<TableRow> reader =
streamSource.createReader(options);
+ reader.start();
+ // Beam does not split storage read api v2 stream
+ assertNull(reader.splitAtFraction(0.5));
+ }
+
@Test
public void testReadFromStreamSource() throws Exception {