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 =