Delta Lake MERGE操作触发INVALID_EXTRACT_BASE_FIELD_TYPE错误的动态字段兼容问题求助
Delta Lake MERGE操作触发INVALID_EXTRACT_BASE_FIELD_TYPE错误的动态字段兼容问题求助
我正在尝试用以下PySpark代码读取数据湖中存储的Delta文件(delta_table),并和包含更新记录的DataFrame(novos_registros)执行MERGE操作:
#5. Build the matching condition for the MERGE condicao_correspondencia = " AND ".join([f"t.{col} = s.{col}" for col in chave_unica]) # 6. Use the MERGE command to update or insert data delta_table = DeltaTable.forPath(spark, delta_path) delta_table.alias("t").merge( novos_registros.alias("s"), condicao_correspondencia # Dynamically generated match condition ).whenMatchedUpdate(set={ **{col: f"s.{col}" for col in novos_registros.columns} # Update all fields of novos_registros }).whenNotMatchedInsert(values={ **{col: f"s.{col}" for col in novos_registros.columns} # Insere todos os campos de novos_registros }).execute()
但现在遇到了字段类型不兼容的问题,报错提示Delta Lake中的0NET_PRICE字段是复杂类型,而DataFrame里的同名字段是字符串类型,具体错误信息如下:
AnalysisException: [INVALID_EXTRACT_BASE_FIELD_TYPE] Can't extract a value from "0NET_PRICE". Need a complex type [STRUCT, ARRAY, MAP] but got "STRING".
看起来DeltaTable.forPath(spark, delta_path)会自动推断Delta表的字段格式,而我的novos_registros DataFrame所有字段都是字符串类型。这是一个通用脚本,我不想在代码里硬编码修复字段格式,希望能实现动态的类型兼容处理。
有没有大佬能指点一下解决办法呀?提前谢谢大家了!
备注:内容来源于stack exchange,提问作者Marcelo Herdy
相关产品推荐
相关产品推荐

