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 b6c8937bbf1 Merge pull request #40274 from
reuvenlax/fix_flaky_schema_test
b6c8937bbf1 is described below
commit b6c8937bbf1d74ebb12cba1c70227f31379ffe69
Author: Reuven Lax <[email protected]>
AuthorDate: Fri Sep 25 08:52:21 2026 -0700
Merge pull request #40274 from reuvenlax/fix_flaky_schema_test
Fix new flaky schema test
---
.../sdk/io/gcp/bigquery/CreateTableHelpers.java | 16 ++++-
.../bigquery/StorageApiWriteUnshardedRecords.java | 78 ++++++++++++++--------
.../bigquery/StorageApiWritesShardedRecords.java | 21 +++++-
.../sdk/io/gcp/testing/FakeDatasetService.java | 38 ++++++++++-
.../sdk/io/gcp/bigquery/BigQueryIOWriteTest.java | 58 ++++++++++++++++
5 files changed, 176 insertions(+), 35 deletions(-)
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/CreateTableHelpers.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/CreateTableHelpers.java
index 7c428917503..52ac77532b1 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/CreateTableHelpers.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/CreateTableHelpers.java
@@ -35,6 +35,7 @@ import
com.google.api.services.bigquery.model.TimePartitioning;
import io.grpc.StatusRuntimeException;
import java.util.Collections;
import java.util.Map;
+import java.util.Optional;
import java.util.Set;
import java.util.concurrent.Callable;
import java.util.concurrent.ConcurrentHashMap;
@@ -48,6 +49,7 @@ import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.Vi
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Strings;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Supplier;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Throwables;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.joda.time.Duration;
@@ -69,13 +71,21 @@ public class CreateTableHelpers {
static void createTableWrapper(Callable<Void> action, Callable<Boolean>
tryCreateTable)
throws Exception {
BackOff backoff =
BackOffAdapter.toGcpBackOff(DEFAULT_BACKOFF_FACTORY.backoff());
- RuntimeException lastException = null;
+ Exception lastException = null;
do {
try {
action.call();
return;
- } catch (ApiException | StatusRuntimeException e) {
- lastException = e;
+ } catch (Exception e) {
+ // The Storage Write library can wrap errors in
UncheckedExecutionException
+ Optional<Throwable> handledCause =
+ Throwables.getCausalChain(e).stream()
+ .filter(
+ cause ->
+ (cause instanceof ApiException || cause instanceof
StatusRuntimeException))
+ .findAny();
+ lastException = (Exception) handledCause.orElseThrow(() -> e);
+
// TODO: Once BigQuery reliably returns a consistent error on table
not found, we should
// only try creating
// the table on that error.
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 33ed612d29c..ecae8dad630 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
@@ -89,6 +89,7 @@ import org.apache.beam.sdk.values.WindowedValues;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Predicates;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Strings;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Throwables;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterators;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists;
@@ -924,13 +925,60 @@ public class
StorageApiWriteUnshardedRecords<DestinationT, ElementT>
quotaError = statusCode.equals(Status.Code.RESOURCE_EXHAUSTED);
}
+ // Schema mismatched exceptions can happen if the table was
recently updated. Since
+ // vortex caches schemas
+ // we might see the new schema before vortex does. In this case,
we simply need to
+ // retry.
+ // Note: ConnectionWorker in google-cloud-bigquerystorage
already converts the gRPC
+ // error via Exceptions.toStorageException(), which strips the
gRPC Status trailers.
+ // Calling Exceptions.toStorageException() a second time on a
StorageException returns
+ // null, so we must check instanceof Exceptions.StorageException
first.
+ Exceptions.@Nullable StorageException storageException = null;
+ if (error instanceof Exceptions.StorageException) {
+ storageException = (Exceptions.StorageException) error;
+ } else if (error != null) {
+ Optional<Throwable> handledCause =
+ Throwables.getCausalChain(error).stream()
+ .filter(cause -> cause instanceof
Exceptions.StorageException)
+ .findAny();
+ if (handledCause.isPresent()) {
+ storageException = (Exceptions.StorageException)
handledCause.get();
+ } else {
+ storageException = Exceptions.toStorageException(error);
+ }
+ }
+ boolean schemaMismatchError =
+ (storageException instanceof
Exceptions.SchemaMismatchedException);
+ if (!schemaMismatchError && error != null) {
+ // There's no special error code for missing required fields,
and that can also
+ // happen due to vortex
+ // being delayed at seeing a new schema. We're forced to parse
the description to
+ // determine that this
+ // has happened.
+ Status status = Status.fromThrowable(error);
+ if (status.getCode() == Status.Code.INVALID_ARGUMENT) {
+ String description = status.getDescription();
+ schemaMismatchError =
+ description != null
+ && (description.contains("incompatible fields")
+ || description.contains(
+ "Input schema has more fields than BigQuery
schema"));
+ }
+ }
+ if (schemaMismatchError) {
+ LOG.info(
+ "Vortex failed stream open due to incompatible fields.
This is likely because the Bigtable "
+ + "schema was recently updated and Vortex hasn't
noticed yet, so retrying. error {}",
+ Preconditions.checkStateNotNull(error).toString());
+ }
+
int allowedRetry;
if (!quotaError) {
// This forces us to close and reopen all gRPC connections to
Storage API on error,
// which empirically fixes random stuckness issues.
invalidateAppendClient(true);
- allowedRetry = 5;
+ allowedRetry = schemaMismatchError ? 35 : 5;
} else {
allowedRetry = 35;
}
@@ -962,34 +1010,6 @@ public class
StorageApiWriteUnshardedRecords<DestinationT, ElementT>
+ failedContext.offset);
}
- // Schema mismatched exceptions can happen if the table was
recently updated. Since
- // vortex caches schemas
- // we might see the new schema before vortex does. In this case,
we simply need to
- // retry.
- Exceptions.@Nullable StorageException storageException =
- (error == null) ? null :
Exceptions.toStorageException(error);
- boolean schemaMismatchError =
- (storageException instanceof
Exceptions.SchemaMismatchedException);
- if (!schemaMismatchError && error != null) {
- // There's no special error code for missing required fields,
and that can also
- // happen due to vortex
- // being delayed at seeing a new schema. We're forced to parse
the description to
- // determine that this
- // has happened.
- Status status = Status.fromThrowable(error);
- if (status.getCode() == Status.Code.INVALID_ARGUMENT) {
- String description = status.getDescription();
- schemaMismatchError =
- description != null &&
description.contains("incompatible fields");
- }
- }
- if (schemaMismatchError) {
- LOG.info(
- "Vortex failed stream open due to incompatible fields.
This is likely because the Bigtable "
- + "schema was recently updated and Vortex hasn't
noticed yet, so retrying. error {}",
- Preconditions.checkStateNotNull(error).toString());
- }
-
boolean hasPersistentErrors =
failedContext.getError() instanceof
Exceptions.StreamFinalizedException
|| statusCode.equals(Status.Code.INVALID_ARGUMENT)
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiWritesShardedRecords.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiWritesShardedRecords.java
index e615182154c..336311e3421 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiWritesShardedRecords.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiWritesShardedRecords.java
@@ -104,6 +104,7 @@ import org.apache.beam.sdk.values.TypeDescriptor;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Predicates;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Strings;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Throwables;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.Cache;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.CacheBuilder;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;
@@ -729,7 +730,20 @@ public class StorageApiWritesShardedRecords<DestinationT
extends @NonNull Object
// vortex caches schemas
// we might see the new schema before vortex does. In this case, we
simply need to
// retry.
- Exceptions.@Nullable StorageException storageException =
Exceptions.toStorageException(error);
+ Exceptions.@Nullable StorageException storageException = null;
+ if (error instanceof Exceptions.StorageException) {
+ storageException = (Exceptions.StorageException) error;
+ } else {
+ Optional<Throwable> handledCause =
+ Throwables.getCausalChain(error).stream()
+ .filter(cause -> cause instanceof Exceptions.StorageException)
+ .findAny();
+ if (handledCause.isPresent()) {
+ storageException = (Exceptions.StorageException) handledCause.get();
+ } else {
+ storageException = Exceptions.toStorageException(error);
+ }
+ }
boolean schemaMismatchError =
(storageException instanceof Exceptions.SchemaMismatchedException);
if (!schemaMismatchError) {
@@ -743,7 +757,10 @@ public class StorageApiWritesShardedRecords<DestinationT
extends @NonNull Object
Status status = Status.fromThrowable(error);
if (status.getCode() == Code.INVALID_ARGUMENT) {
String description = status.getDescription();
- schemaMismatchError = description != null &&
description.contains("incompatible fields");
+ schemaMismatchError =
+ description != null
+ && (description.contains("incompatible fields")
+ || description.contains("Input schema has more fields
than BigQuery schema"));
}
}
if (schemaMismatchError) {
diff --git
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/testing/FakeDatasetService.java
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/testing/FakeDatasetService.java
index 549c2798226..67dc802f118 100644
---
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/testing/FakeDatasetService.java
+++
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/testing/FakeDatasetService.java
@@ -113,8 +113,24 @@ public class FakeDatasetService implements DatasetService,
WriteStreamService, S
.setCode(code)
.setErrorMessage(errorMessage)
.build();
+ int grpcCode = io.grpc.Status.Code.OK.value();
+ if (code == StorageError.StorageErrorCode.SCHEMA_MISMATCH_EXTRA_FIELDS) {
+ grpcCode = io.grpc.Status.Code.INVALID_ARGUMENT.value();
+ } else if (code == StorageError.StorageErrorCode.STREAM_NOT_FOUND) {
+ grpcCode = io.grpc.Status.Code.NOT_FOUND.value();
+ } else if (code == StorageError.StorageErrorCode.STREAM_FINALIZED) {
+ grpcCode = io.grpc.Status.Code.FAILED_PRECONDITION.value();
+ } else if (code == StorageError.StorageErrorCode.OFFSET_OUT_OF_RANGE) {
+ grpcCode = io.grpc.Status.Code.OUT_OF_RANGE.value();
+ } else if (code == StorageError.StorageErrorCode.OFFSET_ALREADY_EXISTS) {
+ grpcCode = io.grpc.Status.Code.ALREADY_EXISTS.value();
+ }
com.google.rpc.Status status =
-
com.google.rpc.Status.newBuilder().addDetails(Any.pack(storageError)).build();
+ com.google.rpc.Status.newBuilder()
+ .setCode(grpcCode)
+ .setMessage(errorMessage)
+ .addDetails(Any.pack(storageError))
+ .build();
return org.apache.beam.sdk.util.Preconditions.checkArgumentNotNull(
Exceptions.toStorageException(status, null));
}
@@ -241,6 +257,8 @@ public class FakeDatasetService implements DatasetService,
WriteStreamService, S
private volatile String appendRowsErrorCode = null;
private volatile String appendRowsErrorDescription = null;
+ private volatile @Nullable StorageError.StorageErrorCode
appendRowsStorageErrorCode = null;
+ private static AtomicInteger appendRowsStorageErrorRemainingCount = new
AtomicInteger(0);
public void setAppendRowsError(Throwable t) {
io.grpc.Status status = io.grpc.Status.fromThrowable(t);
@@ -248,6 +266,13 @@ public class FakeDatasetService implements DatasetService,
WriteStreamService, S
this.appendRowsErrorDescription = status.getDescription();
}
+ public void setAppendRowsStorageError(
+ StorageError.StorageErrorCode code, String description, int
failureCount) {
+ this.appendRowsStorageErrorCode = code;
+ this.appendRowsErrorDescription = description;
+ appendRowsStorageErrorRemainingCount.set(failureCount);
+ }
+
Map<String, List<String>> insertErrors = Maps.newHashMap();
// The counter for the number of insertions performed.
@@ -257,6 +282,7 @@ public class FakeDatasetService implements DatasetService,
WriteStreamService, S
synchronized (FakeDatasetService.class) {
tables = HashBasedTable.create();
insertCount = new AtomicInteger(0);
+ appendRowsStorageErrorRemainingCount = new AtomicInteger(0);
writeStreams = Maps.newHashMap();
FakeJobService.setUp();
}
@@ -811,6 +837,16 @@ public class FakeDatasetService implements DatasetService,
WriteStreamService, S
@Override
public ApiFuture<AppendRowsResponse> appendRows(long offset, ProtoRows
rows)
throws Exception {
+ if (appendRowsStorageErrorCode != null
+ && appendRowsStorageErrorRemainingCount.getAndDecrement() > 0) {
+ return ApiFutures.immediateFailedFuture(
+ getStorageException(
+ streamName,
+ appendRowsStorageErrorCode,
+ appendRowsErrorDescription != null
+ ? appendRowsErrorDescription
+ : "Storage error"));
+ }
if (appendRowsErrorCode != null) {
io.grpc.Status.Code code =
io.grpc.Status.Code.valueOf(appendRowsErrorCode);
io.grpc.Status status =
diff --git
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOWriteTest.java
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOWriteTest.java
index 3fbc4f5aed5..fdfd2699236 100644
---
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOWriteTest.java
+++
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOWriteTest.java
@@ -1723,6 +1723,64 @@ public class BigQueryIOWriteTest implements Serializable
{
p.run().waitUntilFinish();
}
+ @Test
+ public void testStorageApiRetryOnSchemaMismatchedException() throws
Exception {
+ assumeTrue(useStorageApi);
+ assumeTrue(!useStreaming || useStorageApiApproximate);
+
+ Table table =
+ new Table()
+ .setTableReference(
+ new TableReference()
+ .setProjectId("project-id")
+ .setDatasetId("dataset-id")
+ .setTableId("table-id"))
+ .setSchema(
+ new TableSchema()
+ .setFields(
+ ImmutableList.of(
+ new
TableFieldSchema().setName("number").setType("INTEGER"))));
+ fakeDatasetService.createTable(table);
+
+ // Inject transient SchemaMismatchedException (which has
Status.Code.INVALID_ARGUMENT and
+ // stripped gRPC trailers after Exceptions.toStorageException conversion)
for the first 2
+ // appendRows attempts, then succeed.
+ fakeDatasetService.setAppendRowsStorageError(
+ com.google.cloud.bigquery.storage.v1.StorageError.StorageErrorCode
+ .SCHEMA_MISMATCH_EXTRA_FIELDS,
+ "Input schema has more fields than BigQuery schema, extra fields:
'numeric_extra' Entity:
projects/project-id/datasets/dataset-id/tables/table-id/streams/_default",
+ 2);
+
+ List<Integer> elements = Lists.newArrayList(1, 2, 3);
+
+ BigQueryIO.Write<Integer> write =
+ BigQueryIO.<Integer>write()
+ .to("project-id:dataset-id.table-id")
+
.withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_NEVER)
+ .withFormatFunction(
+ (SerializableFunction<Integer, TableRow>)
+ input -> new TableRow().set("number", input))
+ .withSchema(
+ new TableSchema()
+ .setFields(
+ ImmutableList.of(
+ new
TableFieldSchema().setName("number").setType("INTEGER"))))
+ .withTestServices(fakeBqServices)
+ .withoutValidation();
+
+ PCollection<Integer> input =
p.apply(Create.of(elements).withCoder(BigEndianIntegerCoder.of()));
+ input.apply("WriteToBQ", write);
+
+ p.run().waitUntilFinish();
+
+ assertThat(
+ fakeDatasetService.getAllRows("project-id", "dataset-id", "table-id"),
+ containsInAnyOrder(
+ new TableRow().set("number", "1"),
+ new TableRow().set("number", "2"),
+ new TableRow().set("number", "3")));
+ }
+
@Test
public void testStreamingStorageApiWriteWithAutoShardingWithErrorHandling()
throws Exception {
assumeTrue(useStreaming);