使用Apache Beam Dataflow无法将关联数据写入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

