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

Reply via email to