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);