Apache Beam分支合并问题:多BigQuery写入后合并失败
Apache Beam分支合并写入BigQuery日志表问题解决
问题场景
我正在编写一个Apache Beam Pipeline,该Pipeline分为三个分支,每个分支将数据写入BigQuery,之后想要合并为一个分支写入另一个BigQuery表用于日志记录,但无法完成分支合并。
原代码
pipeline_options = PipelineOptions(None) p = beam.Pipeline(options=pipeline_options) ingest_data = ( p | 'Start Pipeline' >> beam.Create([None]) ) p1 = (ingest_data | 'Read from and to date 1' >> beam.ParDo(OutputValueProviderFn('table1')) | 'fetch API data 1' >> beam.ParDo(get_api_data()) | 'write into gbq 1' >> beam.io.gcp.bigquery.WriteToBigQuery(table='proj.dataset.table1', write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, custom_gcs_temp_location='gs://project/temp') ) p2 = (ingest_data | 'Read from and to date 2' >> beam.ParDo(OutputValueProviderFn('table2')) | 'fetch API data 2' >> beam.ParDo(get_api_data()) | 'write into gbq 2' >> beam.io.gcp.bigquery.WriteToBigQuery(table='proj.dataset.table2', write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, custom_gcs_temp_location='gs://proj/temp') ) p3 = (ingest_data | 'Read from and to date 3' >> beam.ParDo(OutputValueProviderFn('table3')) | 'fetch API data 3' >> beam.ParDo(get_api_data()) | 'write into gbq 3' >> beam.io.gcp.bigquery.WriteToBigQuery(table='proj.dataset.table3', write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, custom_gcs_temp_location='gs://proj/temp') ) # 尝试合并分支写入日志表,此处报错 merge = (p1, p2, p3) | 'Write Log' >> beam.io.gcp.bigquery.WriteToBigQuery(table='proj.dataset.table_logging', write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, custom_gcs_temp_location='gs://proj/temp') )
报错信息
AttributeError: Error trying to access nonexistent attribute `0` in write result. Please see __documentation__ for available attributes.
核心需求
不关心三个分支的输出,仅需合并分支以确保之前的三次写入操作已完成后,再写入日志表。
问题根源
WriteToBigQuery返回的是写入结果的元组对象(包含写入行数、错误信息等统计数据),并非可继续处理的PCollection,因此无法直接将这些结果合并后传给下一个WriteToBigQuery操作。
解决方案
通过以下步骤实现分支合并与日志写入:
- 每个分支完成BigQuery写入后,生成一个标识该分支完成的信号(比如包含表名、完成状态的日志记录)
- 使用
beam.Flatten()合并三个分支的信号PCollection - 基于合并后的
PCollection写入日志表,确保只有当三个分支的写入操作全部完成后,才执行日志写入
修改后的完整代码
import datetime import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions # 定义生成日志信号的DoFn class GenerateLogSignal(beam.DoFn): def __init__(self, table_name): self.table_name = table_name def process(self, unused_element): # 生成包含分支完成信息的日志记录 yield { 'table_name': self.table_name, 'completion_time': datetime.datetime.now().isoformat(), 'status': 'SUCCESS' } pipeline_options = PipelineOptions(None) p = beam.Pipeline(options=pipeline_options) ingest_data = ( p | 'Start Pipeline' >> beam.Create([None]) ) # 分支1:写入table1后生成日志信号 p1 = (ingest_data | 'Read from and to date 1' >> beam.ParDo(OutputValueProviderFn('table1')) | 'fetch API data 1' >> beam.ParDo(get_api_data()) | 'write into gbq 1' >> beam.io.gcp.bigquery.WriteToBigQuery( table='proj.dataset.table1', write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, custom_gcs_temp_location='gs://project/temp' ) | 'Generate log signal 1' >> beam.ParDo(GenerateLogSignal('table1')) ) # 分支2:写入table2后生成日志信号 p2 = (ingest_data | 'Read from and to date 2' >> beam.ParDo(OutputValueProviderFn('table2')) | 'fetch API data 2' >> beam.ParDo(get_api_data()) | 'write into gbq 2' >> beam.io.gcp.bigquery.WriteToBigQuery( table='proj.dataset.table2', write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, custom_gcs_temp_location='gs://proj/temp' ) | 'Generate log signal 2' >> beam.ParDo(GenerateLogSignal('table2')) ) # 分支3:写入table3后生成日志信号 p3 = (ingest_data | 'Read from and to date 3' >> beam.ParDo(OutputValueProviderFn('table3')) | 'fetch API data 3' >> beam.ParDo(get_api_data()) | 'write into gbq 3' >> beam.io.gcp.bigquery.WriteToBigQuery( table='proj.dataset.table3', write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, custom_gcs_temp_location='gs://proj/temp' ) | 'Generate log signal 3' >> beam.ParDo(GenerateLogSignal('table3')) ) # 合并三个分支的日志信号 merged_signals = ( (p1, p2, p3) | 'Flatten signals' >> beam.Flatten() ) # 写入日志表 (merged_signals | 'Write Log' >> beam.io.gcp.bigquery.WriteToBigQuery( table='proj.dataset.table_logging', write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, custom_gcs_temp_location='gs://proj/temp' ) ) # 启动Pipeline并等待完成 p.run().wait_until_finish()
关键说明
GenerateLogSignalDoFn会在对应分支的BigQuery写入完成后输出一条日志记录,保证信号生成依赖于写入操作的完成beam.Flatten()负责将三个分支的信号PCollection合并为一个,确保所有分支完成后才会向日志表写入数据- 最终的日志表会收到三条记录,分别对应三个分支的完成状态,也可以根据需求调整日志内容
内容的提问来源于stack exchange,提问作者Francesco Pegoraro
相关产品推荐
相关产品推荐

