如何在Apache Beam ETL流水线中顺序执行读写与表合并任务?
如何在Apache Beam流水线中实现顺序执行逻辑?
问题背景
我在Google Cloud中基于BigQuery编写了Apache Beam代码,包含两个DoFn类:
ReadExcel:读取Cloud Storage中的Excel文件并转换为字典格式MergeTables:执行BigQuery的MERGE操作,对比源表与目标表的哈希值进行更新或插入
需求是先完成读取Excel并写入BigQuery源表的任务,待其完全执行后再执行源表与目标表的合并操作,但当前流水线代码运行报错,无法保证顺序执行。
现有代码
核心类定义
import apache_beam as beam from apache_beam.io.gcp.gcsio import GcsIO from apache_beam.options.pipeline_options import PipelineOptions from apache_beam.io.gcp.gcsfilesystem import GCSFileSystem import argparse import logging import pandas as pd import datetime from google.cloud import bigquery import json from google.cloud import storage from operator import add from functools import reduce from apache_beam import pvalue # 加载配置 with open('config.json', 'r') as config_file: config = json.load(config_file) credintial = config['project'] project_name = credintial['project_name'] dataset_name = credintial['dataset_name'] table_name_source = credintial['table_name_source'] table_name_target = credintial['table_name_target'] department = credintial['department'] file_name = credintial['file_name'] bucket = credintial['bucket'] table_schema = credintial['table_schema'] output_table_name = "{}.{}".format(dataset_name, table_name_source) today = datetime.date.today() today = today.strftime("%Y%m%d") class ReadExcel(beam.DoFn): def process(self, file_path): with GcsIO().open(file_path, "rb") as file: excel_df = pd.read_excel(file, sheet_name='Sheet1') # 语法错误:return应在process方法内部 return [row.fillna('').to_dict() for _, row in excel_df.iterrows()] class MergeTables(beam.DoFn): def __init__(self, project_name): self.project_name = project_name # 语法错误:process方法应在MergeTables类内部 def process(self, element): sql = f""" MERGE `{project_name}.{dataset_name}.{table_name_target}` t USING `{project_name}.{dataset_name}.{table_name_source}` s ON t.purchase_requisition = s.purchase_requisition AND t.pr_item = s.pr_item WHEN MATCHED AND t.hashed_row != s.hashed_row THEN UPDATE SET -- 替换为实际列 t.hashed_row = s.hashed_row WHEN NOT MATCHED THEN INSERT (-- 替换为实际列, hashed_row) VALUES (-- 替换为实际列, s.hashed_row) """ client = bigquery.Client(project=self.project_name) query_job = client.query(sql) query_job.result() # 等待查询完成
流水线运行代码
def run(argv=None): parser = argparse.ArgumentParser() parser.add_argument('--output', dest='output', required=False, help='Output BQ table to write results to.', default=output_table_name) known_args, pipeline_args = parser.parse_known_args(argv) pipeline_options = PipelineOptions(pipeline_args) with beam.Pipeline(options=pipeline_options) as p: gcs = GCSFileSystem(PipelineOptions(pipeline_args)) # 匹配GCS中的文件 match = gcs.match([f'gs://{bucket}/smb-data/{today}/{department}/{file_name}.xlsx']) if match: file_metadata = match[0] excel_data = ( p | 'Read and Write file' >> beam.Create([file_metadata.path]) | 'Read Excel' >> beam.ParDo(ReadExcel()) ) # 将数据写入BigQuery源表 _ = ( excel_data | 'Write to BigQuery' >> beam.io.WriteToBigQuery( known_args.output, schema=table_schema, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE ) ) # 执行表合并(当前与写操作并行执行) _ = ( p | 'Dummy' >> beam.Create([None]) | 'Merge Tables' >> beam.ParDo(MergeTables(project_name)) ) if __name__ == "__main__": logging.getLogger().setLevel(logging.INFO) run()
解决方案
1. 修复代码中的语法错误
首先解决两个缩进问题,否则代码会直接报错:
- 将
ReadExcel类中return语句缩进至process方法内部,与with块同级 - 将
process方法缩进至MergeTables类内部,作为类的成员方法
2. 建立流水线的依赖关系,保证顺序执行
当前流水线中,写BigQuery和合并操作是并行分支(都直接从Pipeline根节点p出发),因此合并操作可能在写表未完成时就执行,导致源表数据不完整。要实现顺序执行,需让合并操作依赖于写表操作的完成:
方法:利用WriteToBigQuery的输出触发合并
WriteToBigQuery会返回一个包含写入状态的PCollection,我们可以将这个PCollection作为合并操作的输入,确保合并操作在所有数据写入完成后才执行。
修改后的流水线代码如下:
def run(argv=None): parser = argparse.ArgumentParser() parser.add_argument('--output', dest='output', required=False, help='Output BQ table to write results to.', default=output_table_name) known_args, pipeline_args = parser.parse_known_args(argv) pipeline_options = PipelineOptions(pipeline_args) with beam.Pipeline(options=pipeline_options) as p: gcs = GCSFileSystem(PipelineOptions(pipeline_args)) # 匹配GCS中的文件 match = gcs.match([f'gs://{bucket}/smb-data/{today}/{department}/{file_name}.xlsx']) if match: file_metadata = match[0] excel_data = ( p | 'Read and Write file' >> beam.Create([file_metadata.path]) | 'Read Excel' >> beam.ParDo(ReadExcel()) ) # 将数据写入BigQuery源表,并获取写入状态的PCollection write_result = ( excel_data | 'Write to BigQuery' >> beam.io.WriteToBigQuery( known_args.output, schema=table_schema, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE ) ) # 基于写操作的结果触发合并:使用CombineGlobally确保所有写入完成后执行一次 _ = ( write_result | 'Wait for Write Completion' >> beam.CombineGlobally(lambda _: None) | 'Merge Tables' >> beam.ParDo(MergeTables(project_name)) ) if __name__ == "__main__": logging.getLogger().setLevel(logging.INFO) run()
关键说明:
CombineGlobally(lambda _: None):将所有写入状态的元素合并为一个单一元素,确保只有当所有数据都写入BigQuery后,才会触发后续的合并操作- 合并操作的输入依赖于写操作的输出,因此流水线会保证写操作完成后再执行合并
3. 额外优化建议
- 在
MergeTables的process方法中添加错误捕获,避免合并失败导致整个流水线崩溃 - 考虑使用
beam.io.gcp.bigquery.BigQueryMerge(如果使用的是较新版本的Beam),它是官方提供的MERGE操作组件,比自定义DoFn更可靠 - 确保
config.json中的权限配置正确,BigQuery客户端需要有足够的权限执行MERGE操作
内容的提问来源于stack exchange,提问作者H. Sayouf
相关产品推荐
相关产品推荐

