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

Apache Beam中WriteToBigQuery用FILE_LOADS+Avro时如何启用死信模式?

问题描述

Apache Beam官方文档建议写入BigQuery时采用死信模式,可通过FailedRows标签从转换输出中获取写入失败的行。

但使用以下代码时:

WriteToBigQuery(
    table=self.bigquery_table_name,
    schema={"fields": self.bigquery_table_schema},
    method=WriteToBigQuery.Method.FILE_LOADS,
    temp_file_format=FileFormat.AVRO,
)

因某条数据Schema不匹配引发异常:

Error message from worker: Traceback (most recent call last):
File 
    "/my_code/apache_beam/io/gcp/bigquery_tools.py", line 1630, 
    in write self._avro_writer.write(row) File "fastavro/_write.pyx", line 647,
    in fastavro._write.Writer.write File "fastavro/_write.pyx", line 376,
    in fastavro._write.write_data File "fastavro/_write.pyx", line 320,
    in fastavro._write.write_record File "fastavro/_write.pyx", line 374,
    in fastavro._write.write_data File "fastavro/_write.pyx", line 283,
    in fastavro._write.write_union ValueError: [] (type <class 'list'>) do not match ['null', 'double'] on field safety_proxy During handling of the above exception, another exception occurred: Traceback (most recent call last): File "apache_beam/runners/common.py", line 1198,
    in apache_beam.runners.common.DoFnRunner.process File "apache_beam/runners/common.py", line 718,
    in apache_beam.runners.common.PerWindowInvoker.invoke_process File "apache_beam/runners/common.py", line 841,
    in apache_beam.runners.common.PerWindowInvoker._invoke_process_per_window File "apache_beam/runners/common.py", line 1334,
    in apache_beam.runners.common._OutputProcessor.process_outputs File "/my_code/apache_beam/io/gcp/bigquery_file_loads.py", line 258,
    in process writer.write(row) File "/my_code/apache_beam/io/gcp/bigquery_tools.py", line 1635,
    in write ex, self._avro_writer.schema, row)).with_traceback(tb) File "/my_code/apache_beam/io/gcp/bigquery_tools.py", line 1630,
    in write self._avro_writer.write(row) File "fastavro/_write.pyx", line 647,
    in fastavro._write.Writer.write File "fastavro/_write.pyx", line 376,
    in fastavro._write.write_data File "fastavro/_write.pyx", line 320,
    in fastavro._write.write_record File "fastavro/_write.pyx", line 374,
    in fastavro._write.write_data File "fastavro/_write.pyx", line 283,
    in fastavro._write.write_union ValueError: Error writing row to Avro: [] (type <class 'list'>) do not match ['null', 'double'] on field safety_proxy Schema: ...

分析可知,Schema不匹配导致fastavro._write.Writer.write失败抛出异常。希望WriteToBigQuery触发死信模式,将格式错误的行作为FailedRows标记的输出返回,而非直接终止管道。

补充示例代码:

from apache_beam import Create
from apache_beam.io.gcp.bigquery import BigQueryWriteFn, WriteToBigQuery
from apache_beam.io.textio import WriteToText

...

valid_rows = [{"some_field_name": i} for i in range(1000000)]
invalid_rows = [{"wrong_field_name": i}]

pcoll = Create(valid_rows + invalid_rows)

# 因1条无效行导致整个管道失败
write_result = (
    pcoll 
    |  WriteToBigQuery(
        table=self.bigquery_table_name,
        schema={
            "fields": [
                {'name': 'some_field_name', 'type': 'INTEGER', 'mode': 'NULLABLE'},
            ]
        },
        method=WriteToBigQuery.Method.FILE_LOADS,
        temp_file_format=FileFormat.AVRO,
    )
)

# 期望WriteToBigQuery部分成功,输出失败行(因少量错误行导致长时管道失败)
(
    write_result[BigQueryWriteFn.FAILED_ROWS] 
    | WriteToText('gs://my_failed_rows/')
)

解决方案

当使用FILE_LOADS模式写入BigQuery时,Avro序列化阶段的错误(比如Schema不匹配)会直接抛出异常终止管道,默认不会将这类错误行路由到FAILED_ROWS输出。要实现死信处理,需要在写入BigQuery之前,先对数据进行Schema校验,提前捕获并分流无效行。

步骤1:实现Schema校验DoFn

编写一个DoFn来校验每行数据是否符合目标BigQuery Schema,将有效行和无效行分别输出:

from apache_beam import DoFn, TaggedOutput
from apache_beam.io.gcp.bigquery_tools import parse_table_schema_from_json
from jsonschema import validate, ValidationError

class ValidateBigQuerySchema(DoFn):
    def __init__(self, schema_json):
        # 将BigQuery Schema转换为JSON Schema格式用于校验
        self.bq_schema = parse_table_schema_from_json(schema_json)
        self.json_schema = self.bq_schema.to_api_repr()

    def process(self, element):
        try:
            # 校验数据是否匹配Schema
            validate(instance=element, schema=self.json_schema)
            yield element
        except ValidationError as e:
            # 将无效行标记为失败行输出,附带错误信息
            yield TaggedOutput('failed_rows', (element, str(e)))

步骤2:在管道中加入校验逻辑

在WriteToBigQuery之前先执行校验,分流有效行和无效行:

# 定义BigQuery Schema(JSON格式)
bq_schema_json = '''{
    "fields": [
        {"name": "some_field_name", "type": "INTEGER", "mode": "NULLABLE"}
    ]
}'''

# 分流有效行和无效行
valid_pcoll, failed_pcoll = (
    pcoll
    | "Validate BigQuery Schema" >> beam.ParDo(ValidateBigQuerySchema(bq_schema_json))
        .with_outputs('failed_rows', main='valid_rows')
)

# 有效行写入BigQuery
valid_pcoll | "Write Valid Rows to BigQuery" >> WriteToBigQuery(
    table=self.bigquery_table_name,
    schema=bq_schema_json,
    method=WriteToBigQuery.Method.FILE_LOADS,
    temp_file_format=FileFormat.AVRO,
)

# 无效行写入死信存储(GCS文本文件)
failed_pcoll | "Write Failed Rows to GCS" >> WriteToText('gs://my_failed_rows/')

补充说明

  • 上述方法在写入前提前拦截无效行,避免Avro序列化阶段抛出异常终止管道。
  • 如果需要捕获BigQuery加载阶段的错误(比如加载时的Schema不匹配、权限问题等),WriteToBigQuery的FAILED_ROWS输出会自动捕获这类错误,无需额外处理。
  • 对于复杂的Schema校验(比如嵌套字段、自定义类型),可以根据需求扩展ValidateBigQuerySchema的逻辑,或者使用Beam提供的BigQuerySchemaValidator工具类(部分版本支持)。

内容的提问来源于stack exchange,提问作者cozos

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 09:45:40