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

stankiewicz 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 9b3c14c9e5b Treat a 404 from BigQuery table and dataset deletes as 
success (#39934)
9b3c14c9e5b is described below

commit 9b3c14c9e5b04f7cc364bf4aebd19aec7a9175dc
Author: Maksym Tymoshyk <[email protected]>
AuthorDate: Tue Sep 29 16:47:08 2026 +0300

    Treat a 404 from BigQuery table and dataset deletes as success (#39934)
    
    BigQueryIO deletes the temporary tables and datasets it creates. A delete
    can succeed at BigQuery and still have its work item fail to commit
    afterwards; the runner then replays the work item and the replayed delete
    gets a 404 because the first attempt already removed the resource.
    
    deleteTable and deleteDataset passed ALWAYS_RETRY, so that 404 was retried
    MAX_RPC_RETRIES times and then thrown. In a streaming job the work item
    retries forever, which stalls a drain.
    
    Both now use DONT_RETRY_NOT_FOUND and swallow an item-not-found error,
    matching how getTable already handles a 404. Every other status keeps its
    existing retry and failure behaviour.
---
 CHANGES.md                                         |  1 +
 .../sdk/io/gcp/bigquery/BigQueryServicesImpl.java  | 63 ++++++++++++++++------
 .../io/gcp/bigquery/BigQueryServicesImplTest.java  | 40 ++++++++++++++
 3 files changed, 87 insertions(+), 17 deletions(-)

diff --git a/CHANGES.md b/CHANGES.md
index 6327e07585a..408f97fad90 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -83,6 +83,7 @@
 * (Go) Fixed a data race on the Prism runner's artifact cache map in 
JobServices ([#32656](https://github.com/apache/beam/issues/32656)).
 * (Java) Fixed the declared schema of the error output of the Kafka write 
SchemaTransform, which wrapped the error schema a second time and did not match 
the rows it emits ([#39760](https://github.com/apache/beam/issues/39760)).
 * (Go) Fixed pubsubio importing a `google.golang.org/genproto` package removed 
in recent releases, which broke builds of Go modules depending on a current 
`genproto` version ([#40018](https://github.com/apache/beam/issues/40018)).
+* (Java) BigQueryIO now treats a 404 when deleting a temporary table or 
dataset as success, so a replayed work item whose earlier attempt already 
deleted it no longer retries forever 
([#24997](https://github.com/apache/beam/issues/24997)).
 * Fixed X (Java/Python) ([#X](https://github.com/apache/beam/issues/X)).
 
 ## Security Fixes
diff --git 
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryServicesImpl.java
 
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryServicesImpl.java
index 5459926f277..7a924c6aa02 100644
--- 
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryServicesImpl.java
+++ 
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryServicesImpl.java
@@ -823,20 +823,35 @@ public class BigQueryServicesImpl implements 
BigQueryServices {
      *
      * <p>Tries executing the RPC for at most {@code MAX_RPC_RETRIES} times 
until it succeeds.
      *
+     * <p>A table that BigQuery reports as not found is treated as deleted 
successfully, since that
+     * is the state the caller asked for.
+     *
      * @throws IOException if it exceeds {@code MAX_RPC_RETRIES} attempts.
      */
     @Override
     public void deleteTable(TableReference tableRef) throws IOException, 
InterruptedException {
-      executeWithRetries(
-          client
-              .tables()
-              .delete(tableRef.getProjectId(), tableRef.getDatasetId(), 
tableRef.getTableId()),
-          String.format(
-              "Unable to delete table: %s, aborting after %d retries.",
-              tableRef.getTableId(), MAX_RPC_RETRIES),
-          Sleeper.DEFAULT,
-          createDefaultBackoff(),
-          ALWAYS_RETRY);
+      try {
+        executeWithRetries(
+            client
+                .tables()
+                .delete(tableRef.getProjectId(), tableRef.getDatasetId(), 
tableRef.getTableId()),
+            String.format(
+                "Unable to delete table: %s, aborting after %d retries.",
+                tableRef.getTableId(), MAX_RPC_RETRIES),
+            Sleeper.DEFAULT,
+            createDefaultBackoff(),
+            DONT_RETRY_NOT_FOUND);
+      } catch (IOException e) {
+        if (!errorExtractor.itemNotFound(e)) {
+          throw e;
+        }
+
+        // a delete can succeed at bigquery and still have its work item fail 
to commit afterwards.
+        // the runner then replays that work item, and the replayed delete 
gets a 404 because the
+        // first attempt already removed the table. failing here would make 
the work item retry
+        // forever, which in a streaming job stalls the drain indefinitely
+        LOG.info("Table {} is already deleted, treating as success.", 
tableRef.getTableId());
+      }
     }
 
     @Override
@@ -969,18 +984,32 @@ public class BigQueryServicesImpl implements 
BigQueryServices {
      *
      * <p>Tries executing the RPC for at most {@code MAX_RPC_RETRIES} times 
until it succeeds.
      *
+     * <p>A dataset that BigQuery reports as not found is treated as deleted 
successfully, since
+     * that is the state the caller asked for.
+     *
      * @throws IOException if it exceeds {@code MAX_RPC_RETRIES} attempts.
      */
     @Override
     public void deleteDataset(String projectId, String datasetId)
         throws IOException, InterruptedException {
-      executeWithRetries(
-          client.datasets().delete(projectId, datasetId),
-          String.format(
-              "Unable to delete table: %s, aborting after %d retries.", 
datasetId, MAX_RPC_RETRIES),
-          Sleeper.DEFAULT,
-          createDefaultBackoff(),
-          ALWAYS_RETRY);
+      try {
+        executeWithRetries(
+            client.datasets().delete(projectId, datasetId),
+            String.format(
+                "Unable to delete table: %s, aborting after %d retries.",
+                datasetId, MAX_RPC_RETRIES),
+            Sleeper.DEFAULT,
+            createDefaultBackoff(),
+            DONT_RETRY_NOT_FOUND);
+      } catch (IOException e) {
+        if (!errorExtractor.itemNotFound(e)) {
+          throw e;
+        }
+
+        // see deleteTable: a replayed work item can find the dataset its own 
earlier attempt
+        // already removed, and treating that 404 as a failure would retry 
forever
+        LOG.info("Dataset {} is already deleted, treating as success.", 
datasetId);
+      }
     }
 
     static class InsertBatchofRowsCallable implements 
Callable<List<InsertErrors>> {
diff --git 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryServicesImplTest.java
 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryServicesImplTest.java
index 8dd02642444..90b2e1fb872 100644
--- 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryServicesImplTest.java
+++ 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryServicesImplTest.java
@@ -569,6 +569,46 @@ public class BigQueryServicesImplTest {
         tableRef, Collections.emptyList(), null, BackOff.STOP_BACKOFF, 
Sleeper.DEFAULT);
   }
 
+  @Test
+  public void testDeleteTableNotFoundSucceeds() throws IOException, 
InterruptedException {
+    setupMockResponses(
+        response -> {
+          when(response.getContentType()).thenReturn(Json.MEDIA_TYPE);
+          when(response.getStatusCode()).thenReturn(404);
+        });
+
+    BigQueryServicesImpl.DatasetServiceImpl datasetService =
+        new BigQueryServicesImpl.DatasetServiceImpl(bigquery, 
PipelineOptionsFactory.create());
+
+    TableReference tableRef =
+        new TableReference()
+            .setProjectId("projectId")
+            .setDatasetId("datasetId")
+            .setTableId("tableId");
+
+    datasetService.deleteTable(tableRef);
+
+    // exactly one response is prepared, so a retry of the 404 would trip the 
Verify inside the mock
+    // request. the assertion is therefore both "did not throw" and "did not 
retry"
+    verifyAllResponsesAreRead();
+  }
+
+  @Test
+  public void testDeleteDatasetNotFoundSucceeds() throws IOException, 
InterruptedException {
+    setupMockResponses(
+        response -> {
+          when(response.getContentType()).thenReturn(Json.MEDIA_TYPE);
+          when(response.getStatusCode()).thenReturn(404);
+        });
+
+    BigQueryServicesImpl.DatasetServiceImpl datasetService =
+        new BigQueryServicesImpl.DatasetServiceImpl(bigquery, 
PipelineOptionsFactory.create());
+
+    datasetService.deleteDataset("projectId", "datasetId");
+
+    verifyAllResponsesAreRead();
+  }
+
   @Test
   public void testIsTableEmptySucceeds() throws Exception {
     TableReference tableRef =

Reply via email to