You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Apache Beam InsertRetryPolicy对部分BigQuery写入错误不生效问题咨询

BigQuery Dynamic Destinations: Unhandled Exceptions Bypassing 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 BigQueryServicesImpl lines 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:

  1. 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 a GoogleJsonResponseException instead of returning error details in the response.
  2. Beam’s error handling gap: In BigQueryServicesImpl, the retry policy logic only runs when futures.get(i).get() returns a response with InsertErrors. If get() 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 catch GoogleJsonResponseException.
  • 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 InsertErrors objects so they trigger the never() 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 08:27:43