Spark中如何将含新列的DataFrame合并至Hudi表
解决Hudi MERGE INTO列不兼容问题
问题原因
你遇到的报错是因为Hudi的MERGE INTO要求源DataFrame(modification)和目标表(origin)的列结构必须兼容——origin表没有c列,而modification有,直接合并就会触发列不匹配的异常。要保留新列c并完成合并,得先扩展origin表的结构,再执行合并操作。
具体解决步骤
1. 确保Hudi表开启Schema Evolution
如果你的origin表创建时没开启Schema Evolution,先通过配置开启(可在Spark SQL会话中设置,或修改表属性):
SET hoodie.schema.evolution.enable=true; SET hoodie.schema.on.read.enable=true;
2. 给origin表添加新列c
用ALTER TABLE语句扩展表结构,根据modification中c列的实际类型指定类型(示例为字符串类型):
ALTER TABLE origin ADD COLUMN c STRING;
3. 执行MERGE INTO合并操作
现在两张表的列结构兼容了,可执行合并逻辑:
MERGE INTO origin o USING modification m ON o.id = m.id WHEN MATCHED THEN UPDATE SET o.a = m.a, o.b = m.b, o.c = m.c WHEN NOT MATCHED THEN INSERT (id, a, b, c) VALUES (m.id, m.a, m.b, m.c)
补充:用Spark DataFrame API实现的方式
如果习惯用代码操作,也可以这样处理:
// 1. 开启Schema Evolution配置 spark.conf.set("hoodie.schema.evolution.enable", "true") spark.conf.set("hoodie.schema.on.read.enable", "true") // 2. 加载origin表,添加c列(默认值设为null)并更新表结构 val originDF = spark.read.format("hudi").load("path/to/origin") val originWithC = originDF.withColumn("c", lit(null).cast(StringType)) originWithC.write.format("hudi") .option("hoodie.table.name", "origin") .option("hoodie.datasource.write.recordkey.field", "id") .option("hoodie.datasource.write.precombine.field", "id") // 根据实际业务设置预合并字段 .mode("overwrite") .save("path/to/origin") // 3. 执行upsert合并modification数据 modificationDF.write.format("hudi") .option("hoodie.table.name", "origin") .option("hoodie.datasource.write.operation", "upsert") .option("hoodie.datasource.write.recordkey.field", "id") .option("hoodie.datasource.write.precombine.field", "id") .mode("append") .save("path/to/origin")
内容的提问来源于stack exchange,提问作者Rinze
相关产品推荐
相关产品推荐

