使用Cloud Dataflow向BigQuery导入CSV遇Term字段匹配错误求助
Cloud Dataflow导入BigQuery字段匹配错误的解决方法
你使用Cloud Dataflow将Cloud Storage中的CSV数据导入BigQuery时遇到字段匹配错误,核心原因是BigQuery表的schema定义了Term字段,但CSV和你的Dataflow代码里对应的字段是T,两者不匹配导致错误。
你的CSV格式:
T,WeekCodeOption,TimeWidthSEQ,StationCode,RealBroadcastMinutes,StandardTargetFlag,TargetInfoCode,Rate,ExcludeDates "{'TFrom': '20230213', 'TTo': '20230213'}",11,1,1,1440,"['1', '1', '1', '1']","['101', '102', '103', '104']","{'101': ['3.2'], '102': ['1.9'], '103': ['1.6'], '104': ['2.6']}",
你的代码中,convert_dict函数返回的字典使用"T"作为字段名,但BigQuery表的schema定义的是Term:STRING,这是矛盾的核心。
方案一:修改Dataflow代码中的字段名
直接在convert_dict函数里,把返回字典中的"T"键改为"Term",让它和BigQuery的schema字段匹配,同时建议将数字类型字段转换为对应类型,避免类型不匹配:
修改后的convert_dict函数:
def convert_dict(element): now_utc = datetime.now(timezone.utc) now = datetime.now(TZ_JST) # 可选:如果需要解析T字段的JSON内容,可以取消注释下面的代码并处理 # term = json.loads(element[0].replace("'", '"')) return { "Term": element[0], # 此处将"T"改为"Term" "WeekCodeOption": int(element[1]), "TimeWidthSEQ": int(element[2]), "StationCode": int(element[3]), "RealBroadcastMinutes": int(element[4]), "StandardTargetFlag": element[5], "TargetInfoCode": element[6], "Rate": element[7], "ExcludeDates": element[8], "created_time_ts_utc": now_utc.strftime("%Y-%m-%d %H:%M:%S.%f"), "created_time_dt_jst": now.strftime("%Y-%m-%d %H:%M:%S.%f"), }
另外,当前代码会把CSV的表头行也当成数据处理,需要添加跳过表头的逻辑:
raw_datas = p | "Read from GCS" >> beam.io.ReadFromText(gcs_uri) # 跳过CSV表头行 filtered_datas = raw_datas | "Skip header" >> beam.Filter(lambda line: not line.startswith("T,")) tran_datas = ( filtered_datas | "transform csv" >> beam.Map(parse_file) | "transform dict" >> beam.Map(convert_dict) )
方案二:修改BigQuery表的schema
如果不想调整代码,可以修改BigQuery表的schema,把Term:STRING改为T:STRING,让schema和CSV及代码中的字段名一致。修改后的schema字符串:
schema = 'T:STRING, WeekCodeOption:INTEGER, TimeWidthSEQ:INTEGER, StationCode:INTEGER, RealBroadcastMinutes:INTEGER, StandardTargetFlag:STRING, TargetInfoCode:STRING, Rate:STRING, ExcludeDates:STRING'
注意:因为你设置了create_disposition=BigQueryDisposition.CREATE_NEVER,需要手动在BigQuery控制台修改已存在表的schema;或者改为CREATE_IF_NEEDED,让Dataflow自动使用修改后的schema创建表。
内容的提问来源于stack exchange,提问作者サワヤン
相关产品推荐
相关产品推荐

