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

使用Dataflow+CoGroupByKey关联数据写入BigQuery时报错

问题描述

以peakid为关联键,通过CoGroupByKey关联expeditions-2010s.csv和peaks.csv两个文件,写入BigQuery时触发以下错误:

RuntimeError: BigQuery job beam_bq_job_LOAD_AUTOMATIC_JOB_NAME_LOAD_STEP_88_215864ba592a2e01f0c4e2157cc60c47_86e3562707f348c29b2a030cb6ed7ded failed. Error Result: <ErrorProto location:'gs://bucket-name/input/temp/bq_load/ededcfb43cda4d16934011481e2fd774/project_name.dataset.expeditions/9fe30f70-8473-44bc-86d5-20dfdf59f502' message: '读取数据时出错,错误消息:JSON表遇到过多错误,已放弃。行数:1;错误数:1。请查看errors[]集合获取详细信息。文件:gs://bucket-name/input/temp/bq_load/ededcfb43cda4d16934011481e2fd774/project_name.dataset.expeditions/9fe30f70-8473-44bc-86d5-20dfdf59f502' reason: 'invalid'> [while running 'Write To BigQuery/BigQueryBatchFileLoads/WaitForDestinationLoadJobs']

原始代码
def read_csv_pd_input1(readable_file):
    import json
    import pandas as pd   
    import csv
    import io
    gcs_file = beam.io.filesystems.FileSystems.open(readable_file)
    csv_dict = csv.DictReader(io.TextIOWrapper(gcs_file))
    df = pd.DataFrame(csv_dict)
    df = df[['peakid', 'bcdate', 'smtdate']]
    
    a = df.set_index('peakid')[['bcdate', 'smtdate']].apply(tuple,1).to_dict()
    a = tuple(a.items())
    
    # result: only column name   
    # a = df.agg(lambda x: (x.values)).apply(tuple)

    # result: only value but not as expected    
    # a = [tuple(x) for x in df.values]
    # a = tuple(a)
    return a

def read_csv_pd_input3(readable_file):
    import json
    import pandas as pd   
    import csv
    import io
    gcs_file = beam.io.filesystems.FileSystems.open(readable_file)
    csv_dict = csv.DictReader(io.TextIOWrapper(gcs_file))
    df = pd.DataFrame(csv_dict)
    df = df[['peakid', 'pkname', 'heightm']] 
    
    a = df.set_index('peakid')[['pkname', 'heightm']].apply(tuple,1).to_dict()
    a = tuple(a.items())
    
    return a


def run(argv=None):
    import apache_beam as beam
    import io

    parser = argparse.ArgumentParser()
    parser.add_argument(
        '--input',
        dest='input',
        required=False,
        help='Input file to read. This can be a local file or '
        'a file in a Google Storage Bucket.',
        default='gs://bucket-name/input/expeditions- 2010s.csv')
    
    parser.add_argument(
        '--input3',
        dest='input3',
        required=False,
        help='Input_p3 file to read. This can be a local file or '
        'a file in a Google Storage Bucket.',
        default='gs://bucket-name/input/peaks.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
         | 'Read From GCS input1' >> beam.Create([known_args.input])
         | 'Pair each employee with key p1' >> beam.FlatMap(read_csv_pd_input1)
         # | beam.Map(print)
        
    )
    input_p3 = (
        p
         | 'Read From GCS input3' >> beam.Create([known_args.input3])
         | 'Pair each employee with key p3' >> beam.FlatMap(read_csv_pd_input3)
    )
    # CoGroupByKey: relational join of 2 or more key/values PCollection. It also accept dictionary of key value
    output = (
        {'input_p1': input_p1, 'input_p3': input_p3} 
        | 'Join' >> beam.CoGroupByKey()
        | 'Write To BigQuery' >> beam.io.gcp.bigquery.WriteToBigQuery(
           table='project_name:dataset.expeditions',
           schema='peakid:STRING,bcdate:DATE,pkname:STRING,heightm:INTEGER',
           method='FILE_LOADS',
           custom_gcs_temp_location='gs://bucket-name/input/temp',
           create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
           write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE)    
    )
    p.run().wait_until_finish()
    # runner = DataflowRunner()
    # runner.run_pipeline(p, options=options)

if __name__ == '__main__':
    logging.getLogger().setLevel(logging.INFO)
    run()
问题分析
  1. 输出结构不匹配BigQuery Schema:CoGroupByKey的输出格式为(peakid, {'input_p1': [(bcdate, smtdate), ...], 'input_p3': [(pkname, heightm), ...]}),是嵌套结构,但BigQuery需要扁平的键值对,直接写入会导致JSON解析失败。
  2. CSV读取逻辑冗余易出错:用Pandas手动读取并转换为字典元组的方式复杂,容易出现数据结构错误,不如Beam原生的ReadFromCsv可靠。
  3. 数据类型未正确转换:原始代码中heightm是字符串类型,未转为整数;bcdate若不符合YYYY-MM-DD格式,会触发BigQuery类型校验错误。
修复方案

1. 简化CSV读取逻辑

用Beam原生ReadFromCsv直接读取,指定列名和类型,避免手动处理:

# 读取探险数据,输出(peakid, bcdate)
def read_expeditions(file_path):
    return (
        beam.io.ReadFromCsv(
            file_path,
            columns=['peakid', 'bcdate', 'smtdate'],
            skip_header_lines=1
        )
        | beam.Map(lambda row: (row['peakid'], row['bcdate']))
    )

# 读取山峰数据,输出(peakid, (pkname, heightm))
def read_peaks(file_path):
    return (
        beam.io.ReadFromCsv(
            file_path,
            columns=['peakid', 'pkname', 'heightm'],
            skip_header_lines=1
        )
        | beam.Map(lambda row: (row['peakid'], (row['pkname'], int(row['heightm']) if row['heightm'].strip() else None)))
    )

2. 转换关联结果结构

添加函数将嵌套的关联结果展开为BigQuery需要的扁平格式,处理一对多或缺失数据的情况:

def flatten_joined_data(element):
    peakid, grouped_data = element
    expeditions = grouped_data.get('input_p1', [])
    peaks = grouped_data.get('input_p3', [])
    
    # 假设一个peakid对应一条山峰数据、多条探险数据,遍历展开
    for exp in expeditions:
        bcdate = exp if exp else None
        pkname = peaks[0][0] if peaks else None
        heightm = peaks[0][1] if peaks else None
        yield {
            'peakid': peakid,
            'bcdate': bcdate,
            'pkname': pkname,
            'heightm': heightm
        }

3. 修改Pipeline主逻辑

替换原始读取和转换步骤,确保数据结构符合BigQuery要求:

def run(argv=None):
    import apache_beam as beam
    import argparse
    import logging

    parser = argparse.ArgumentParser()
    parser.add_argument(
        '--input',
        dest='input',
        required=False,
        help='输入文件路径(本地或GCS)',
        default='gs://bucket-name/input/expeditions- 2010s.csv')
    
    parser.add_argument(
        '--input3',
        dest='input3',
        required=False,
        help='山峰数据文件路径(本地或GCS)',
        default='gs://bucket-name/input/peaks.csv')
     
    known_args, pipeline_args = parser.parse_known_args(argv)

    pipeline_options = beam.PipelineOptions(pipeline_args)

    with beam.Pipeline(options=pipeline_options) as p:
        # 读取并转换探险数据
        input_p1 = (
            p
            | '读取探险数据' >> beam.io.ReadFromCsv(
                known_args.input,
                columns=['peakid', 'bcdate', 'smtdate'],
                skip_header_lines=1
            )
            | '转换探险数据为键值对' >> beam.Map(lambda row: (row['peakid'], row['bcdate']))
        )

        # 读取并转换山峰数据
        input_p3 = (
            p
            | '读取山峰数据' >> beam.io.ReadFromCsv(
                known_args.input3,
                columns=['peakid', 'pkname', 'heightm'],
                skip_header_lines=1
            )
            | '转换山峰数据为键值对' >> beam.Map(lambda row: (row['peakid'], (row['pkname'], int(row['heightm']) if row['heightm'].strip() else None)))
        )

        # 关联数据并转换为扁平结构,写入BigQuery
        output = (
            {'input_p1': input_p1, 'input_p3': input_p3} 
            | '关联数据' >> beam.CoGroupByKey()
            | '展开关联结果' >> beam.FlatMap(flatten_joined_data)
            | '写入BigQuery' >> beam.io.gcp.bigquery.WriteToBigQuery(
               table='project_name:dataset.expeditions',
               schema='peakid:STRING,bcdate:DATE,pkname:STRING,heightm:INTEGER',
               method='FILE_LOADS',
               custom_gcs_temp_location='gs://bucket-name/input/temp',
               create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
               write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE)    
        )

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

额外注意事项

  • 确保bcdate列格式为YYYY-MM-DD,若原始数据格式不符,需添加日期转换逻辑。
  • 处理空值:将空字符串转为None,BigQuery会自动处理为NULL。
  • 可在写入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.20 03:40:29