This is an automated email from the ASF dual-hosted git repository.

reuvenlax 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 dd89aa1e4e3 Merge pull request #40269 from 
reuvenlax/schema_update_fixups
dd89aa1e4e3 is described below

commit dd89aa1e4e380dc73bc219a97757c3e8c3ec2d61
Author: Reuven Lax <[email protected]>
AuthorDate: Fri Sep 25 08:59:00 2026 -0700

    Merge pull request #40269 from reuvenlax/schema_update_fixups
    
    Some followup fixups to the schema-update code
---
 .../beam/sdk/io/gcp/bigquery/BigQueryOptions.java  |  2 +-
 .../io/gcp/bigquery/StorageApiWritePayload.java    |  3 +-
 .../bigquery/StorageApiWriteUnshardedRecords.java  |  4 +-
 .../sdk/io/gcp/bigquery/AppendRowsPacketTest.java  |  6 +--
 .../SchemaChangeDetectorHelperBufferingTest.java   | 17 +++++++-
 .../bigquery/SchemaChangeDetectorHelperTest.java   | 46 +++++++++++-----------
 .../StorageApiSchemaMismatchDrainTest.java         |  8 ++--
 .../bigquery/StorageApiSinkSchemaUpdateITBase.java |  4 +-
 8 files changed, 52 insertions(+), 38 deletions(-)

diff --git 
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryOptions.java
 
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryOptions.java
index da8526bd994..1826f4a1fd2 100644
--- 
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryOptions.java
+++ 
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryOptions.java
@@ -133,7 +133,7 @@ public interface BigQueryOptions
   @Description(
       "When using the STORAGE_API_AT_LEAST_ONCE write method with multiplexing 
(ie. useStorageApiConnectionPool=true), "
           + "this option sets the maximum number of connections each pool 
creates. This is on a per worker, per region basis. "
-          + "If writing to many dynamic destinations (>20) and experiencing 
performance issues or seeing append operations competing"
+          + "If writing to many dynamic destinations (>20) and experiencing 
performance issues or seeing append operations competing "
           + "for streams, consider increasing this value.")
   @Default.Integer(20)
   Integer getMaxConnectionPoolConnections();
diff --git 
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiWritePayload.java
 
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiWritePayload.java
index 6e54c3da404..d60f8135597 100644
--- 
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiWritePayload.java
+++ 
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiWritePayload.java
@@ -87,8 +87,7 @@ public abstract class StorageApiWritePayload {
       @Nullable Instant timestamp,
       @Nullable byte[] unknownFieldsPayload,
       @Nullable byte[] failsafeTableRowPayload,
-      @Nullable byte[] schemaHash)
-      throws IOException {
+      @Nullable byte[] schemaHash) {
     return new AutoValue_StorageApiWritePayload.Builder()
         .setPayload(payload)
         .setTimestamp(timestamp)
diff --git 
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiWriteUnshardedRecords.java
 
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiWriteUnshardedRecords.java
index ecae8dad630..8ea8d75ed17 100644
--- 
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiWriteUnshardedRecords.java
+++ 
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiWriteUnshardedRecords.java
@@ -529,7 +529,7 @@ public class StorageApiWriteUnshardedRecords<DestinationT, 
ElementT>
       }
 
       AppendClientInfo generateClient(@Nullable TableSchema updatedSchema) 
throws Exception {
-        SchemaAndDescriptor schemaAndDescriptor = 
getCurrentTableSchema(streamName, updatedSchema);
+        SchemaAndDescriptor schemaAndDescriptor = 
getCurrentTableSchema(updatedSchema);
 
         AtomicReference<AppendClientInfo> appendClientInfo =
             new AtomicReference<>(
@@ -566,7 +566,7 @@ public class StorageApiWriteUnshardedRecords<DestinationT, 
ElementT>
         }
       }
 
-      SchemaAndDescriptor getCurrentTableSchema(String stream, @Nullable 
TableSchema updatedSchema)
+      private SchemaAndDescriptor getCurrentTableSchema(@Nullable TableSchema 
updatedSchema)
           throws Exception {
         if (updatedSchema != null) {
           return new SchemaAndDescriptor(
diff --git 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/AppendRowsPacketTest.java
 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/AppendRowsPacketTest.java
index 1b1e64dd029..562ed122898 100644
--- 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/AppendRowsPacketTest.java
+++ 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/AppendRowsPacketTest.java
@@ -80,11 +80,11 @@ public class AppendRowsPacketTest {
   }
 
   private static Instant timestampFor(int i) {
-    return new Instant(1_000 + i);
+    return Instant.ofEpochMilli(1_000 + i);
   }
 
   private static Instant deadlineFor(int i) {
-    return new Instant(2_000 + i);
+    return Instant.ofEpochMilli(2_000 + i);
   }
 
   private static StoragePayloadWithDeadline payloadFor(int i) {
@@ -251,7 +251,7 @@ public class AppendRowsPacketTest {
             Iterators.peekingIterator(inputs.iterator()),
             Long.MAX_VALUE,
             helper,
-            new Instant(5_000),
+            Instant.ofEpochMilli(5_000),
             appendClientInfo,
             e -> false);
 
diff --git 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/SchemaChangeDetectorHelperBufferingTest.java
 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/SchemaChangeDetectorHelperBufferingTest.java
index 7f11f05ac4a..83b101ede39 100644
--- 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/SchemaChangeDetectorHelperBufferingTest.java
+++ 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/SchemaChangeDetectorHelperBufferingTest.java
@@ -27,6 +27,9 @@ import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
 
 import com.google.api.services.bigquery.model.TableRow;
+import com.google.cloud.bigquery.storage.v1.TableFieldSchema;
+import com.google.cloud.bigquery.storage.v1.TableSchema;
+import com.google.protobuf.DescriptorProtos;
 import java.util.List;
 import org.apache.beam.sdk.metrics.Counter;
 import org.apache.beam.sdk.metrics.Metrics;
@@ -77,14 +80,24 @@ public class SchemaChangeDetectorHelperBufferingTest {
 
   @Before
   @SuppressWarnings("unchecked")
-  public void setUp() {
+  public void setUp() throws Exception {
     bufferedBag = new FakeBagState<>();
     currentTimerValue = new FakeValueState<>();
     minPendingTimestamp = new FakeValueState<>();
     retryTimer = new FakeTimer(NOW);
     tableDestination = new TableDestination("project-id:dataset-id.table", 
null);
     counter = Metrics.counter(SchemaChangeDetectorHelperBufferingTest.class, 
"failedRows");
-    appendClientInfo = mock(AppendClientInfo.class);
+    TableSchema tableSchema =
+        TableSchema.newBuilder()
+            .addFields(
+                TableFieldSchema.newBuilder()
+                    .setName("name")
+                    .setType(TableFieldSchema.Type.STRING)
+                    .build())
+            .build();
+    DescriptorProtos.DescriptorProto descriptor =
+        TableRowToStorageApiProto.descriptorSchemaFromTableSchema(tableSchema, 
true, false);
+    appendClientInfo = AppendClientInfo.of(tableSchema, descriptor, client -> 
{});
     failedRowsReceiver = mock(DoFn.OutputReceiver.class);
   }
 
diff --git 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/SchemaChangeDetectorHelperTest.java
 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/SchemaChangeDetectorHelperTest.java
index 526259e16e6..486c63e0539 100644
--- 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/SchemaChangeDetectorHelperTest.java
+++ 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/SchemaChangeDetectorHelperTest.java
@@ -21,8 +21,6 @@ import static org.junit.Assert.assertArrayEquals;
 import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertFalse;
 import static org.junit.Assert.assertTrue;
-import static org.mockito.ArgumentMatchers.any;
-import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.when;
 
@@ -48,14 +46,24 @@ public class SchemaChangeDetectorHelperTest {
   private TableReference tableReference;
   private BigQueryServices.WriteStreamService mockWriteStreamService;
   private BigQueryServices.StreamAppendClient mockStreamAppendClient;
-  private AppendClientInfo mockAppendClientInfo;
+  private AppendClientInfo appendClientInfo;
 
   @Before
-  public void setUp() {
+  public void setUp() throws Exception {
     tableReference = new 
TableReference().setProjectId("p").setDatasetId("d").setTableId("t");
     mockWriteStreamService = mock(BigQueryServices.WriteStreamService.class);
     mockStreamAppendClient = mock(BigQueryServices.StreamAppendClient.class);
-    mockAppendClientInfo = mock(AppendClientInfo.class);
+    TableSchema tableSchema =
+        TableSchema.newBuilder()
+            .addFields(
+                TableFieldSchema.newBuilder()
+                    .setName("foo")
+                    .setType(TableFieldSchema.Type.STRING)
+                    .build())
+            .build();
+    DescriptorProtos.DescriptorProto descriptor =
+        TableRowToStorageApiProto.descriptorSchemaFromTableSchema(tableSchema, 
true, false);
+    appendClientInfo = AppendClientInfo.of(tableSchema, descriptor, client -> 
{});
   }
 
   @Test
@@ -146,7 +154,7 @@ public class SchemaChangeDetectorHelperTest {
         StorageApiWritePayload.of(new byte[] {1, 2, 3}, new 
TableRow().set("foo", "bar"), null);
 
     SchemaChangeDetectorHelper.MergePayloadResult result =
-        helper.getMergedPayload(payload, Instant.now(), null, 
mockAppendClientInfo);
+        helper.getMergedPayload(payload, Instant.now(), null, 
appendClientInfo);
 
     assertEquals(SchemaChangeDetectorHelper.MergePayloadResult.Kind.MERGED, 
result.getKind());
     assertArrayEquals(new byte[] {1, 2, 3}, result.getMerged().toByteArray());
@@ -159,7 +167,7 @@ public class SchemaChangeDetectorHelperTest {
     StorageApiWritePayload payload = StorageApiWritePayload.of(new byte[] {1, 
2, 3}, null, null);
 
     SchemaChangeDetectorHelper.MergePayloadResult result =
-        helper.getMergedPayload(payload, Instant.now(), null, 
mockAppendClientInfo);
+        helper.getMergedPayload(payload, Instant.now(), null, 
appendClientInfo);
 
     assertEquals(SchemaChangeDetectorHelper.MergePayloadResult.Kind.MERGED, 
result.getKind());
     assertArrayEquals(new byte[] {1, 2, 3}, result.getMerged().toByteArray());
@@ -173,37 +181,31 @@ public class SchemaChangeDetectorHelperTest {
     StorageApiWritePayload payload =
         StorageApiWritePayload.of(new byte[] {1, 2, 3}, unknownFields, null);
 
-    ByteString mergedBytes = ByteString.copyFrom(new byte[] {4, 5, 6});
-    when(mockAppendClientInfo.mergeNewFields(any(ByteString.class), 
eq(unknownFields), eq(false)))
-        .thenReturn(mergedBytes);
+    ByteString expectedMerged =
+        appendClientInfo.mergeNewFields(
+            ByteString.copyFrom(new byte[] {1, 2, 3}), unknownFields, false);
 
     SchemaChangeDetectorHelper.MergePayloadResult result =
-        helper.getMergedPayload(payload, Instant.now(), null, 
mockAppendClientInfo);
+        helper.getMergedPayload(payload, Instant.now(), null, 
appendClientInfo);
 
     assertEquals(SchemaChangeDetectorHelper.MergePayloadResult.Kind.MERGED, 
result.getKind());
-    assertArrayEquals(new byte[] {4, 5, 6}, result.getMerged().toByteArray());
+    assertArrayEquals(expectedMerged.toByteArray(), 
result.getMerged().toByteArray());
   }
 
   @Test
   public void testGetMergedPayload_autoUpdateTrue_mergeFailure() throws 
Exception {
     SchemaChangeDetectorHelper helper =
         new SchemaChangeDetectorHelper(true, false, tableReference, false);
-    TableRow unknownFields = new TableRow().set("foo", "bar");
-    StorageApiWritePayload payload =
-        StorageApiWritePayload.of(new byte[] {1, 2, 3}, unknownFields, null);
-
-    when(mockAppendClientInfo.mergeNewFields(any(ByteString.class), 
eq(unknownFields), eq(false)))
-        .thenThrow(new 
TableRowToStorageApiProto.SchemaDoesntMatchException("conversion error"));
+    TableRow unknownFields = new TableRow().set("unknown_col", "bar");
+    StorageApiWritePayload payload = StorageApiWritePayload.of(new byte[0], 
unknownFields, null);
     TableRow expectedFailsafe = new TableRow().set("failsafe", "true");
 
     SchemaChangeDetectorHelper.MergePayloadResult result =
-        helper.getMergedPayload(payload, Instant.now(), expectedFailsafe, 
mockAppendClientInfo);
+        helper.getMergedPayload(payload, Instant.now(), expectedFailsafe, 
appendClientInfo);
 
     assertEquals(SchemaChangeDetectorHelper.MergePayloadResult.Kind.FAILED, 
result.getKind());
     TimestampedValue<BigQueryStorageApiInsertError> failed = 
result.getFailed();
-    assertEquals(
-        
"org.apache.beam.sdk.io.gcp.bigquery.TableRowToStorageApiProto$SchemaDoesntMatchException:
 conversion error",
-        failed.getValue().getErrorMessage());
+    assertTrue(failed.getValue().getErrorMessage().contains("unknown_col"));
     assertEquals(expectedFailsafe, failed.getValue().getRow());
   }
 
diff --git 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSchemaMismatchDrainTest.java
 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSchemaMismatchDrainTest.java
index 8ff950ebfb0..eeaef4073b5 100644
--- 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSchemaMismatchDrainTest.java
+++ 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSchemaMismatchDrainTest.java
@@ -202,7 +202,7 @@ public class StorageApiSchemaMismatchDrainTest implements 
Serializable {
     bqOptions.setStorageApiMismatchDrainRetryTimeMilliSec(1000);
 
     TestStream.Builder<Long> testStream =
-        TestStream.create(VarLongCoder.of()).advanceWatermarkTo(new 
Instant(0));
+        
TestStream.create(VarLongCoder.of()).advanceWatermarkTo(Instant.ofEpochMilli(0));
     for (long i = 0; i < NUM_ROWS; i++) {
       testStream = testStream.addElements(i);
     }
@@ -370,7 +370,7 @@ public class StorageApiSchemaMismatchDrainTest implements 
Serializable {
     bqOptions.setStorageApiMismatchRetryTimeMilliSec(500);
 
     TestStream.Builder<Long> testStream =
-        TestStream.create(VarLongCoder.of()).advanceWatermarkTo(new 
Instant(0));
+        
TestStream.create(VarLongCoder.of()).advanceWatermarkTo(Instant.ofEpochMilli(0));
     for (long i = 0; i < NUM_ROWS; i++) {
       testStream = testStream.addElements(i);
     }
@@ -442,7 +442,7 @@ public class StorageApiSchemaMismatchDrainTest implements 
Serializable {
     bqOptions.setStorageApiMismatchDrainRetryTimeMilliSec(5_000);
 
     TestStream.Builder<Long> testStream =
-        TestStream.create(VarLongCoder.of()).advanceWatermarkTo(new 
Instant(0));
+        
TestStream.create(VarLongCoder.of()).advanceWatermarkTo(Instant.ofEpochMilli(0));
     for (long i = 0; i < NUM_ROWS; i++) {
       testStream = testStream.addElements(i);
     }
@@ -593,7 +593,7 @@ public class StorageApiSchemaMismatchDrainTest implements 
Serializable {
             true);
 
     TestStream.Builder<Long> testStream =
-        TestStream.create(VarLongCoder.of()).advanceWatermarkTo(new 
Instant(0));
+        
TestStream.create(VarLongCoder.of()).advanceWatermarkTo(Instant.ofEpochMilli(0));
     for (long i = 0; i < NUM_ROWS; i++) {
       testStream = testStream.addElements(i);
     }
diff --git 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkSchemaUpdateITBase.java
 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkSchemaUpdateITBase.java
index 18a94cf85b9..5c27fcd283f 100644
--- 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkSchemaUpdateITBase.java
+++ 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiSinkSchemaUpdateITBase.java
@@ -544,7 +544,7 @@ abstract class StorageApiSinkSchemaUpdateITBase {
     // set up and build pipeline.
     // Rows are emitted as fast as possible; any wall-clock delay that the 
test needs is inserted
     // by UpdateSchemaDoFn around the schema change itself.
-    Instant start = new Instant(0);
+    Instant start = Instant.ofEpochMilli(0);
     Duration interval = Duration.millis(1);
     Duration stop = Duration.millis(TOTAL_N - 1);
     Function<Instant, Long> getIdFromInstant =
@@ -831,7 +831,7 @@ abstract class StorageApiSinkSchemaUpdateITBase {
 
     int numRows = TOTAL_N;
     // set up and build pipeline
-    Instant start = new Instant(0);
+    Instant start = Instant.ofEpochMilli(0);
     // We give a healthy waiting period between each element to give Storage 
API streams a chance to
     // recognize the new schema. Apply on relevant tests.
     Duration interval = Duration.millis(1);

Reply via email to