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

Apache Beam Dataflow中CoGroupByKey结果写入BigQuery报错排查

排查与解决:CoGroupByKey关联数据写入BigQuery失败问题

问题背景

两个CSV文件以peakid为关联键,通过Apache Beam的CoGroupByKey完成关联后,写入BigQuery时触发数据格式错误。

数据形态

expeditions-2010s(关联前)

  • ('TKRG', {'bcdate': ['3/6/10'], 'smtdate': ['3/12/10']})
  • ('AMPG', {'bcdate': ['4/5/10'], 'smtdate': ['']})
  • ('AMAD', {'bcdate': ['4/5/10'], 'smtdate': ['4/21/10']})
  • ...

peak(关联前)

  • ('ACHN', {'pkname': ['Aichyn'], 'heightm': ['6055']})
  • ('AGLE', {'pkname': ['Agole East'], 'heightm': ['6675']})
  • ('AMAD', {'pkname': ['Ama Dablam'], 'heightm': ['6814']})
  • ...

CoGroupByKey关联后输出

  • ('ACHN', ([{'bcdate': [''], 'smtdate': ['9/25/15']}, ...], [{'pkname': ['Aichyn'], 'heightm': ['6055']}]))
  • ('AGLE', ([], [{'pkname': ['Agole East'], 'heightm': ['6675']}]))
  • ('AMAD', ([{'bcdate': ['4/5/10'], ...}, ...], [{'pkname': ['Ama Dablam'], ...}]))
  • ...

错误信息

BigQuery job beam_bq_job_LOAD_AUTOMATIC_JOB_NAME_LOAD_STEP_460_215864ba592a2e01f0c4e2157cc60c47_bc7734af2ebb4a53a0e268bbe6c40824 failed. Error Result: <ErrorProto
location: 'gs://bucket-name/input/temp/bq_load/ece048e1a1ed41b987210a5c4b5e2c52/project-name.dataset.expeditions/cdcdbb44-2e25-4f4a-a792-34382d828244'
message: 'Error while reading data, error message: JSON table encountered too many errors, giving up. Rows: 1; errors: 1. Please look into the errors[] collection for more details. File: gs://bucket-name/input/temp/bq_load/ece048e1a1ed41b987210a5c4b5e2c52/project-name.dataset.expeditions/cdcdbb44-2e25-4f4a-a792-34382d828244'
reason: 'invalid'> [while running 'Write To BigQuery/BigQueryBatchFileLoads/WaitForDestinationLoadJobs']

原始代码片段

读取并转换expeditions数据

input_p1 = (
        p
         | 'Read From GCS input1' >> beam.Create([known_args.input1])
         | 'Parse csv file p1' >> beam.FlatMap(read_csv_file)
         | 'Tuple p1' >> beam.Map(lambda e: (e["peakid"], {'bcdate': [e["bcdate"]], 'smtdate':[e["smtdate"]]}))
    )

读取并转换peak数据

input_p2 = (
        p
         | 'Read From GCS input2' >> beam.Create([known_args.input])
         | 'Parse csv file p2' >> beam.FlatMap(read_csv_file)
         | 'Tuple p2' >> beam.Map(lambda e: (e["peakid"], {'pkname': [e["pkname"]], 'heightm':[e["heightm"]]}))
    )

关联并写入BigQuery

output = (
        (input_p1, input_p2)
        | 'Join' >> beam.CoGroupByKey()
        # | beam.Map(print)
        | '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://dtnhu_test_dataflow_v1/input/temp',
           create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
           write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE)    
    )

问题分析

  1. 数据结构不匹配:CoGroupByKey输出是(peakid, (expeditions_list, peak_list))的嵌套结构,但BigQuery要求扁平的键值对,无法直接解析列表类型的字段。
  2. 日期格式不兼容:BigQuery的DATE类型要求YYYY-MM-DD格式,但当前数据是MM/DD/YY格式,空字符串也无法直接转为DATE类型。
  3. 字段类型不匹配:heightm在BigQuery中定义为INTEGER,但当前是字符串列表,且未处理空值情况。
  4. 一对多关联未展开:单个peak对应多个expedition的情况需要展开为多行记录,而非保留列表。

解决方案

  1. 展开嵌套结构:将CoGroupByKey的输出展开,每个expedition条目与对应的peak信息组合成单行记录;无expedition的peak可选择过滤或填充默认值。
  2. 转换日期格式:将MM/DD/YY格式转为YYYY-MM-DD,空字符串转为None(BigQuery会识别为NULL)。
  3. 类型转换:将heightm的字符串转为整数,处理空值。
  4. 数据清洗:确保每个字段符合BigQuery的类型要求。

修正后的代码

import datetime

def parse_date(date_str):
    """转换日期格式为BigQuery兼容的YYYY-MM-DD,空值返回None"""
    if not date_str.strip():
        return None
    try:
        dt = datetime.datetime.strptime(date_str, '%m/%d/%y')
        return dt.strftime('%Y-%m-%d')
    except ValueError:
        return None

def process_joined_data(element):
    """展开CoGroupByKey的结果,转换字段格式"""
    peakid, (expeditions, peaks) = element
    # 处理peak信息(每个peakid对应一条peak记录)
    peak_info = peaks[0] if peaks else {'pkname': None, 'heightm': None}
    pkname = peak_info['pkname'][0] if peak_info.get('pkname') else None
    heightm = int(peak_info['heightm'][0]) if (peak_info.get('heightm') and peak_info['heightm'][0].strip()) else None
    
    # 展开expedition记录,每个expedition生成一条数据
    if expeditions:
        for exp in expeditions:
            bcdate = parse_date(exp['bcdate'][0]) if exp.get('bcdate') else None
            yield {
                'peakid': peakid,
                'bcdate': bcdate,
                'pkname': pkname,
                'heightm': heightm
            }
    else:
        # 无expedition的情况,生成一条空bcdate的记录(可根据需求过滤)
        yield {
            'peakid': peakid,
            'bcdate': None,
            'pkname': pkname,
            'heightm': heightm
        }

# 原始输入处理部分不变
input_p1 = (
        p
         | 'Read From GCS input1' >> beam.Create([known_args.input1])
         | 'Parse csv file p1' >> beam.FlatMap(read_csv_file)
         | 'Tuple p1' >> beam.Map(lambda e: (e["peakid"], {'bcdate': [e["bcdate"]], 'smtdate':[e["smtdate"]]}))
    )

input_p2 = (
        p
         | 'Read From GCS input2' >> beam.Create([known_args.input])
         | 'Parse csv file p2' >> beam.FlatMap(read_csv_file)
         | 'Tuple p2' >> beam.Map(lambda e: (e["peakid"], {'pkname': [e["pkname"]], 'heightm':[e["heightm"]]}))
    )

# 关联后处理并写入BigQuery
output = (
        (input_p1, input_p2)
        | 'Join' >> beam.CoGroupByKey()
        | 'Process Joined Data' >> beam.FlatMap(process_joined_data)
        | '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://dtnhu_test_dataflow_v1/input/temp',
           create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
           write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE)    
    )

额外说明

  • 如果不需要保留无expedition的peak记录,可在process_joined_data函数中去掉else分支。
  • 可根据实际数据情况调整日期解析的异常处理逻辑,避免因无效日期导致任务失败。
  • 若read_csv_file返回的字段本身不是列表,可简化Tuple p1和Tuple p2中的映射逻辑,直接存储单个值而非列表。

内容的提问来源于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.19 19:31:02