diff --git a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/TableRowToStorageApiProto.java b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/TableRowToStorageApiProto.java index 371c77c5f9c1..ba72bb8682fd 100644 --- a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/TableRowToStorageApiProto.java +++ b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/TableRowToStorageApiProto.java @@ -1639,9 +1639,9 @@ public static ByteString mergeNewFields( null, null, collectedExceptions); - if (!collectedExceptions.isEmpty()) { - return null; - } + } + if (!collectedExceptions.isEmpty()) { + return null; } } else if (schemaInformation.getType() == TableFieldSchema.Type.TIMESTAMP && schemaInformation.getTimestampPrecision() == PICOSECOND_PRECISION) { 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 ef390f5a8601..b9ab25ec5810 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 @@ -36,6 +36,7 @@ import org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO.Write.CreateDisposition; import org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO.Write.WriteDisposition; import org.apache.beam.sdk.io.gcp.testing.BigqueryClient; +import org.apache.beam.sdk.options.StreamingOptions; import org.apache.beam.sdk.state.StateSpec; import org.apache.beam.sdk.state.StateSpecs; import org.apache.beam.sdk.state.ValueState; @@ -187,6 +188,7 @@ public void testDataTriggeredSchemaUpgradeAtLeastOnce() throws Exception { private void runTest(Write.Method method) throws Exception { Pipeline p = Pipeline.create(TestPipeline.testingPipelineOptions()); p.getOptions().as(BigQueryOptions.class).setSchemaUpgradeBufferingShards(1); + p.getOptions().as(StreamingOptions.class).setStreaming(true); TableSchema baseSchema = new TableSchema()