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

如何在Apache Beam单个Pipeline中依次执行BigQuery读写操作

单个Apache Beam Pipeline内实现串行执行的方案

核心思路

由于WriteToBigQuery没有返回可迭代的PCollection,无法直接作为下游依赖,需要通过显式依赖控制让第二部分等待第一部分写入完成后再执行。

方案一:使用Wait变换强制依赖

通过创建空PCollection,让第二部分等待第一部分的写入完成信号,确保串行执行。

修改后的完整代码:

from apache_beam.transforms.util import Wait

p = beam.Pipeline(options=opts)

# 第一部分:写入BigQuery
part_1 = (
    p
    | "F1: Read data 1" >> beam.io.ReadFromText(entrada, skip_header_lines=True)
    | "F1: Transform 1" >> beam.Map(format_date)
    | "F1: Transform 2" >> beam.Map(make_row)
    | "F1: Write into BQ" >> WriteToBigQuery(
        output_table,
        schema=table_schema,
        write_disposition=BigQueryDisposition.WRITE_APPEND,
        create_disposition=BigQueryDisposition.CREATE_IF_NEEDED,
        additional_bq_parameters={
            "timePartitioning": {"type": "DAY"},
            "clustering": {"fields": ["programcode"]},
        },
        custom_gcs_temp_location=temp_location
    )
)

# 创建空信号PCollection,用于绑定依赖
signal = p | "Create execution signal" >> beam.Create([])

# 第二部分:等待第一部分完成后读取BigQuery
part_2 = (
    signal
    | "Wait for part1 write finish" >> Wait(part_1)
    | "F2: Read from BigQuery" >> beam.io.ReadFromBigQuery(
        query=query_raw,
        use_standard_sql=True,
        gcs_location=temp_location,
        project=project_id
    )
)

result = p.run()
result.wait_until_finish()

方案二:复用第一部分的中间结果作为依赖

直接绑定第一部分写入前的PCollection作为依赖,Beam执行引擎会保证写入完成后再启动第二部分:

p = beam.Pipeline(options=opts)

# 提取第一部分的中间处理结果
processed_data = (
    p
    | "F1: Read data 1" >> beam.io.ReadFromText(entrada, skip_header_lines=True)
    | "F1: Transform 1" >> beam.Map(format_date)
    | "F1: Transform 2" >> beam.Map(make_row)
)

# 第一部分:写入BigQuery
part_1 = processed_data | "F1: Write into BQ" >> WriteToBigQuery(
    output_table,
    schema=table_schema,
    write_disposition=BigQueryDisposition.WRITE_APPEND,
    create_disposition=BigQueryDisposition.CREATE_IF_NEEDED,
    additional_bq_parameters={
        "timePartitioning": {"type": "DAY"},
        "clustering": {"fields": ["programcode"]},
    },
    custom_gcs_temp_location=temp_location
)

# 第二部分:绑定中间结果依赖,确保写入完成后读取
part_2 = (
    processed_data
    | "Bind dependency to part1" >> beam.Map(lambda x: x)  # 无意义变换,仅用于绑定执行顺序
    | "F2: Read from BigQuery" >> beam.io.ReadFromBigQuery(
        query=query_raw,
        use_standard_sql=True,
        gcs_location=temp_location,
        project=project_id
    )
)

result = p.run()
result.wait_until_finish()

关键说明

  • Wait变换是最直接的串行控制方式,会严格等待上游所有操作完成后再启动下游。
  • 方案二中的无意义Map变换,目的是让第二部分的读取操作绑定到第一部分的处理流程,触发Beam的依赖调度逻辑。
  • 这种串行控制会牺牲部分并行性,需根据业务场景的优先级权衡使用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 23:22:47