Apache Beam InsertRetryPolicy对部分BigQuery写入错误不生效问题咨询
InsertRetryPolicy The Problem
I’ve been using BigQuery IO’s Dynamic Destinations API to write data to tables, and I set InsertRetryPolicy.never() to avoid retries on errors. But here’s the catch:
- Errors like schema mismatches work as expected—they go through the retry policy logic in
BigQueryServicesImpllines 751-771, and don’t trigger retries. - However, two specific error types are bypassing this entirely: they throw uncaught exceptions from the underlying BigQuery client when executing
List<TableDataInsertAllResponse.InsertErrors> errors = futures.get(i).get()(line 751). This causes the entire bundle to fail, and since it’s a streaming job, the bundle retries infinitely, blocking the whole pipeline.
The unhandled exceptions I’m seeing look like this:
404 Table Not Found
{ "code" : 404, "errors" : [ { "domain" : "global", "message" : "Not found: Table account:dataset.paymentfraudresultsv1_uk_2018", "reason" : "notFound" } ] }
400 Invalid Partition Boundary
{ "code" : 400, "errors" : [ { "domain" : "global", "message" : "The destination table's partition whatever_uk_2017$20171020 is outside the allowed bounds. You can only stream to partitions within 31 days in the past and 16 days in the future relative to the current date.", "reason" : "invalid" } ] }
I’m not sure why some errors come back as structured responses while others throw exceptions, but this is breaking pipeline stability badly.
Why This Happens
Looking into the BigQuery Java client and Beam’s BigQuery IO implementation, here’s what’s going on:
- Client-side error categorization: The BigQuery client treats certain errors as "pre-request failures" or "unrecoverable" and throws exceptions immediately, instead of wrapping them in a
TableDataInsertAllResponse. For example, missing tables (404) or invalid partitions (400) are considered configuration issues that need fixing upfront, not data-level errors that can be handled during writes. So it throws aGoogleJsonResponseExceptioninstead of returning error details in the response. - Beam’s error handling gap: In
BigQueryServicesImpl, the retry policy logic only runs whenfutures.get(i).get()returns a response withInsertErrors. Ifget()throws an exception, the code skips the error handling branch entirely and propagates the exception up, triggering a bundle retry.
Fixes & Workarounds
Here are a few approaches to resolve this:
1. Pre-Validate Tables and Partitions
Add a pre-processing step before writing to BigQuery:
- For dynamically generated tables, call the BigQuery API to check if the table exists first. If it doesn’t, either create it (if your workflow allows) or filter out the data destined for that table.
- For partitioned tables, validate that the partition date falls within BigQuery’s allowed window (31 days past to 16 days future). If not, route that data to an error handling path or discard it.
This is the cleanest solution because it prevents these problematic requests from being sent to BigQuery in the first place.
2. Extend BigQuery IO to Catch These Exceptions
If pre-validation isn’t feasible, you can modify Beam’s BigQuery IO implementation to catch these exceptions and route them through the retry policy logic:
- Wrap the
futures.get(i).get()call in a try-catch block to catchGoogleJsonResponseException. - Parse the exception’s status code and error reason to identify the problematic errors (404 notFound, 400 invalid with partition bounds).
- Convert these exceptions into
InsertErrorsobjects so they trigger thenever()retry policy instead of causing a bundle retry.
Here’s a rough code snippet showing how you might adjust the logic in BigQueryServicesImpl:
try { List<TableDataInsertAllResponse.InsertErrors> errors = futures.get(i).get(); // Existing error handling logic } catch (ExecutionException e) { if (e.getCause() instanceof GoogleJsonResponseException) { GoogleJsonResponseException jsonEx = (GoogleJsonResponseException) e.getCause(); int statusCode = jsonEx.getStatusCode(); ErrorInfo errorInfo = jsonEx.getDetails().getErrors().get(0); if ((statusCode == 404 && "notFound".equals(errorInfo.getReason())) || (statusCode == 400 && "invalid".equals(errorInfo.getReason()) && jsonEx.getMessage().contains("partition"))) { // Convert exception to InsertErrors to trigger retry policy TableDataInsertAllResponse.InsertErrors insertError = new TableDataInsertAllResponse.InsertErrors(); // Populate error details from the exception insertError.setErrors(Collections.singletonList(errorInfo)); // Add to your errors list and continue processing errorsCollection.add(insertError); continue; } } // Re-throw other exceptions as usual throw e; }
3. Use a Dead-Letter Queue
Configure a dead-letter queue (DLQ) in your Beam pipeline to route unprocessable data to a separate storage location (like GCS or a dedicated BigQuery table) instead of letting the bundle retry. You can use BigQueryIO.Write.withFailedInsertRetryPolicy() alongside withDeadLetterQueue() to capture these error records and keep the pipeline running smoothly.
Wrap-Up
This issue stems from a mismatch between how the BigQuery client categorizes errors and how Beam’s BigQuery IO handles them. Pre-validation is the most straightforward fix, but extending the IO logic or using a DLQ are solid alternatives if your workflow can’t support upfront checks.
内容的提问来源于stack exchange,提问作者Luke De Feo

