开启Delta自动合并Schema后读取表出现找不到新列的异常问题
Delta表追加新列后查询报错的解决方法
问题重现步骤
- 开启Delta自动合并Schema配置:
spark.conf.set("spark.databricks.delta.schema.autoMerge.enabled","true")
- 创建初始DataFrame:
data = [("data0",10,"2023-06-26 08:16:18.036", 18.5), ("data1",24,"2023-06-26 01:19:18.036", 1.5), ("data2",25,"2023-06-26 09:20:18.036", 18.5)] schema = StructType([StructField("col1",StringType(),True), StructField("col2",IntegerType(),True), StructField("col3",StringType(),True), StructField("col4",FloatType(),True)]) df_mirror = spark.createDataFrame(data=data,schema=schema) df_mirror = df_mirror.withColumn("col3",col("col3").cast("timestamp"))
- 创建并写入
schema_test.landing2表:
df_landing = df_mirror.withColumn("ins", lit("NOMDB"))\ .withColumn("insert_date", current_timestamp() - expr("INTERVAL 24 HOURS") )\ .withColumn("update_date",current_timestamp() - expr("INTERVAL 24 HOURS")) df_landing.write.mode("append").saveAsTable("schema_test.landing2")
此时表数据结构为:
+-----+----+-----------------------+----+--------+-----------------------+-----------------------+ |col1 |col2|col3 |col4|ins |insert_date |update_date | +-----+----+-----------------------+----+--------+-----------------------+-----------------------+ |data0|10 |2023-06-26 08:16:18.036|18.5|NOMDB |2023-07-05 20:53:03.526|2023-07-05 20:53:03.526| |data1|24 |2023-06-26 01:19:18.036|1.5 |NOMDB |2023-07-05 20:53:03.526|2023-07-05 20:53:03.526| |data2|25 |2023-06-26 09:20:18.036|18.5|NOMDB |2023-07-05 20:53:03.526|2023-07-05 20:53:03.526| +-----+----+-----------------------+----+--------+-----------------------+-----------------------+
- 创建含新增列
col5的DataFrame并追加写入:
+-----+----+-----------------------+----+----+--------+-----------------------+-----------------------+ |col1 |col2|col3 |col4|col5|ins |insert_date |update_date | +-----+----+-----------------------+----+----+--------+-----------------------+-----------------------+ |data8|10 |2023-06-26 08:16:18.036|18.5|null|NOMDB |2023-07-05 20:56:16.891|2023-07-05 20:56:16.891| +-----+----+-----------------------+----+----+--------+-----------------------+-----------------------+
执行追加代码:
df_landing.write.mode("append").saveAsTable("schema_test.landing2")
- 查询表时触发错误:
SELECT * FROM schema_test.landing2
错误信息:
java.lang.IllegalStateException: Couldn't find col5#25652 in [col1#25645,col2#25646,col3#25647,col4#25648,instance#25649,insert_date#25650,update_date#25651]
原因分析
spark.databricks.delta.schema.autoMerge.enabled仅在Merge操作时生效,普通append写入不会触发自动合并Schema。- 追加含新列的数据时,Delta表的元数据未更新,导致查询时新旧数据文件的Schema不兼容,引发报错。
解决方法
方法1:用Merge操作替代Append(推荐)
利用Delta Lake的Merge功能,既可以追加数据,又能自动合并Schema:
from delta.tables import DeltaTable delta_table = DeltaTable.forName(spark, "schema_test.landing2") # 若仅需追加无需匹配更新,设置匹配条件为"1=0"即可 delta_table.alias("target").merge( df_landing.alias("source"), "target.col1 = source.col1" # 根据业务逻辑设置匹配条件 ).whenNotMatchedInsertAll().execute()
方法2:手动合并Schema后重写表
读取现有表,合并新列Schema后重新写入:
# 读取现有表 existing_df = spark.read.table("schema_test.landing2") # 合并新旧DataFrame的Schema combined_df = existing_df.unionByName(df_landing, allowMissingColumns=True) # 覆盖写入更新表元数据 combined_df.write.mode("overwrite").saveAsTable("schema_test.landing2")
方法3:用Overwrite模式并开启mergeSchema(谨慎使用)
若允许覆盖原有数据,可开启mergeSchema参数写入:
df_landing.write.mode("overwrite").option("mergeSchema", "true").saveAsTable("schema_test.landing2")
注:此方法会丢失原有数据,仅适用于数据可重放的场景。
内容的提问来源于stack exchange,提问作者BryC
相关产品推荐
相关产品推荐

