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
相关产品推荐
相关产品推荐

