使用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()
- 输出结构不匹配BigQuery Schema:CoGroupByKey的输出格式为
(peakid, {'input_p1': [(bcdate, smtdate), ...], 'input_p3': [(pkname, heightm), ...]}),是嵌套结构,但BigQuery需要扁平的键值对,直接写入会导致JSON解析失败。 - CSV读取逻辑冗余易出错:用Pandas手动读取并转换为字典元组的方式复杂,容易出现数据结构错误,不如Beam原生的
ReadFromCsv可靠。 - 数据类型未正确转换:原始代码中
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

