Apache Beam Dataflow中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) )
问题分析
- 数据结构不匹配:
CoGroupByKey输出是(peakid, (expeditions_list, peak_list))的嵌套结构,但BigQuery要求扁平的键值对,无法直接解析列表类型的字段。 - 日期格式不兼容:BigQuery的
DATE类型要求YYYY-MM-DD格式,但当前数据是MM/DD/YY格式,空字符串也无法直接转为DATE类型。 - 字段类型不匹配:
heightm在BigQuery中定义为INTEGER,但当前是字符串列表,且未处理空值情况。 - 一对多关联未展开:单个peak对应多个expedition的情况需要展开为多行记录,而非保留列表。
解决方案
- 展开嵌套结构:将
CoGroupByKey的输出展开,每个expedition条目与对应的peak信息组合成单行记录;无expedition的peak可选择过滤或填充默认值。 - 转换日期格式:将
MM/DD/YY格式转为YYYY-MM-DD,空字符串转为None(BigQuery会识别为NULL)。 - 类型转换:将
heightm的字符串转为整数,处理空值。 - 数据清洗:确保每个字段符合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

