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

使用Apache Beam Dataflow无法将关联数据写入BigQuery

Apache Beam Dataflow关联CSV写入BigQuery失败排查

问题场景

有两个CSV文件:expeditions- 2010s.csv和peaks.csv,通过peak_id字段关联,使用Apache Beam在Dataflow中实现关联后写入BigQuery,但写入失败。

关联代码

def read_csv_file(readable_file):
    import apache_beam as beam
    import csv
    import io
    import datetime
    # 打开GCS文件通道
    gcs_file = beam.io.filesystems.FileSystems.open(readable_file)

    # 读取为CSV字典格式
    csv_dict = csv.DictReader(io.TextIOWrapper(gcs_file))

    for row in csv_dict:
        yield (row)

def run(argv=None):
    import apache_beam as beam
    import io
    parser = argparse.ArgumentParser()
    parser.add_argument(
        '--input',
        dest='input',
        required=False,
        help='输入文件路径,支持本地或GCS存储桶',
        default='gs://bucket/folder/peaks.csv')
    
    parser.add_argument(
        '--input1',
        dest='input1',
        required=False,
        help='输入文件路径,支持本地或GCS存储桶',
        default='gs://bucket/folder/expeditions- 2010s.csv')

     
    known_args, pipeline_args = parser.parse_known_args(argv)

    pipeline_options = PipelineOptions(pipeline_args)

    p = beam.Pipeline(options=PipelineOptions(pipeline_args))
    input_p1 = (
        p
         | '读取GCS输入1' >> beam.Create([known_args.input1])
         | '解析CSV文件p1' >> beam.FlatMap(read_csv_file)
         | '转换为键值对p1' >> beam.Map(lambda e: (e["peakid"], {'peakid': e["peakid"], 'bcdate': e["bcdate"], 'smtdate':e["smtdate"]}))
    )
    input_p2 = (
        p
         | '读取GCS输入2' >> beam.Create([known_args.input])
         | '解析CSV文件p2' >> beam.FlatMap(read_csv_file)
         | '转换为键值对p2' >> beam.Map(lambda e: (e["peakid"], {'peakid': e["peakid"], 'pkname': e["pkname"], 'heightm':e["heightm"]})) 
    )
    # CoGroupByKey:关联多个键值对PCollection
    output = (
        (input_p1, input_p2)
        | '关联数据' >> beam.CoGroupByKey()
        | '合并为最终字典' >> beam.Map(lambda el: to_final_dict(el[1]))
        # | beam.Map(print)
        | '写入BigQuery' >> beam.io.gcp.bigquery.WriteToBigQuery(
           table='project:dataset.expeditions',
           method='FILE_LOADS',
           custom_gcs_temp_location='gs://bucket/folder/temp',
           create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
           write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE)    
    )
    p.run().wait_until_finish()
    
def to_final_dict(list_tuple_of_tuple):
        result = {}
        for list_tuple in list_tuple_of_tuple:
            for el in list_tuple:
                result.update(el)
        return result

if __name__ == '__main__':
    logging.getLogger().setLevel(logging.INFO)
    run()

写入前输出示例

  • {'peakid': 'TKRG', 'bcdate': '4/24/10', 'smtdate': '5/5/10', 'pkname': 'Takargo', 'heightm': '6771'}
  • {'peakid': 'AMPG', 'bcdate': '4/5/10', 'smtdate': '', 'pkname': 'Amphu Gyabjen', 'heightm': '5630'}
  • {'peakid': 'AMAD', 'bcdate': '1/27/20', 'smtdate': '2/2/20', 'pkname': 'Ama Dablam', 'heightm': '6814'}
  • {'peakid': 'ANN1', 'bcdate': '3/27/19', 'smtdate': '4/23/19', 'pkname': 'Annapurna I', 'heightm': '8091'}
  • ...

报错信息(中文翻译)

RuntimeError: BigQuery任务beam_bq_job_LOAD_AUTOMATIC_JOB_NAME_LOAD_STEP_602_215864ba592a2e01f0c4e2157cc60c47_51de5de53b58409da70f699c833c4db5执行失败。错误详情:<ErrorProto
location: 'gs://bucket/folder/temp/bq_load/4bbfc44d750c4af5ab376b2e3c3dedbd/project.dataset.expeditions/25905e46-db76-49f0-9b98-7d77131e3e0d'
message: '读取数据时出错,错误信息:JSON表遇到过多错误,已终止。处理行数:3;错误数:1。请查看errors[]集合获取详细信息。文件:gs://bucket/folder/temp/bq_load/4bbfc44d750c4af5ab376b2e3c3dedbd/project.dataset.expeditions/25905e46-db76-49f0-9b98-7d77131e3e0d'
reason: 'invalid'> [运行阶段:'Write To BigQuery/BigQueryBatchFileLoads/WaitForDestinationLoadJobs']

排查方向与解决方法

1. 获取详细错误信息

进入Google Cloud Console的BigQuery页面,找到报错中的加载任务ID,查看errors字段的具体内容,这是定位问题的核心。常见具体错误包括:

  • 数据类型不匹配:比如heightm在BigQuery中是数值类型,但代码传递的是字符串;
  • 日期格式不兼容:BigQuery默认支持YYYY-MM-DD,当前日期是MM/DD/YY格式;
  • 空值处理问题:空字符串''不符合BigQuery日期/数值字段的要求;
  • 字段名不匹配:输出字段名与BigQuery表字段名大小写或名称不一致。

2. 针对性修复示例

(1)转换日期格式

在数据处理阶段将日期转为BigQuery兼容格式:

from datetime import datetime

def parse_date(date_str):
    if not date_str:
        return None
    # 把MM/DD/YY转为YYYY-MM-DD
    return datetime.strptime(date_str, '%m/%d/%y').strftime('%Y-%m-%d')

# 修改input_p1的Map步骤
| '转换为键值对p1' >> beam.Map(lambda e: (e["peakid"], {
    'peakid': e["peakid"], 
    'bcdate': parse_date(e["bcdate"]), 
    'smtdate': parse_date(e["smtdate"])
}))

(2)转换数值类型

将heightm从字符串转为数值:

# 修改input_p2的Map步骤
| '转换为键值对p2' >> beam.Map(lambda e: (e["peakid"], {
    'peakid': e["peakid"], 
    'pkname': e["pkname"], 
    'heightm': float(e["heightm"]) if e["heightm"] else None
}))

(3)显式定义表结构

避免BigQuery自动推断结构出错,在写入时指定表结构:

table_schema = {
    'fields': [
        {'name': 'peakid', 'type': 'STRING', 'mode': 'REQUIRED'},
        {'name': 'bcdate', 'type': 'DATE', 'mode': 'NULLABLE'},
        {'name': 'smtdate', 'type': 'DATE', 'mode': 'NULLABLE'},
        {'name': 'pkname', 'type': 'STRING', 'mode': 'NULLABLE'},
        {'name': 'heightm', 'type': 'FLOAT64', 'mode': 'NULLABLE'}
    ]
}

# 修改WriteToBigQuery参数
| '写入BigQuery' >> beam.io.gcp.bigquery.WriteToBigQuery(
    table='project:dataset.expeditions',
    schema=table_schema,
    method='FILE_LOADS',
    custom_gcs_temp_location='gs://bucket/folder/temp',
    create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
    write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE)

3. 调试建议

暂时注释写入BigQuery的步骤,添加beam.Map(print)将输出数据打印到日志,检查是否存在异常格式的数据。


内容的提问来源于stack exchange,提问作者Nhu Dao

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 21:15:40