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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 10:24:54