如何在Apache Spark中合并两个DataFrame并以第二个DataFrame的值覆盖
在Apache Spark中合并DataFrame并覆盖重复键的方案
刚好碰到过类似的需求,这里给你几个实用的实现方法,都能达到你要的效果——合并两个DataFrame,当键重复时用第二个DataFrame的值覆盖第一个。
方法1:Union + 去重(通用简洁方案)
这个方法的思路是先把两个DataFrame合并,再根据键去重,保留最后出现的行(也就是第二个DataFrame里的重复行)。要注意顺序:先放第一个DF,再放第二个DF,这样去重时会自动保留后面的行。
代码示例:
val r1 = Seq((1, "A1_1"), (2, "A2_1"), (3, "A3_1"), (4, "A4_1")).toDF("c1","c2") val r2 = Seq((3, "A3_2"), (4, "A4_2"), (5, "A5_2"), (6, "A6_2")).toDF("c1","c2") // 合并后按c1去重,保留最后出现的行 val mergedDF = r1.union(r2).dropDuplicates("c1").orderBy("c1") mergedDF.show()
输出结果和你期望的完全一致:
+---+----+ | c1| c2| +---+----+ | 1|A1_1| | 2|A2_1| | 3|A3_2| | 4|A4_2| | 5|A5_2| | 6|A6_2| +---+----+
方法2:Full Outer Join + Coalesce(字段灵活处理方案)
如果你的DataFrame有多个字段,需要针对每个字段单独设置取值优先级,这个方法会更灵活。用full_outer关联两个DF,再用coalesce函数优先选择第二个DF的字段值,最后整理成目标结构。
代码示例:
import org.apache.spark.sql.functions._ val r1 = Seq((1, "A1_1"), (2, "A2_1"), (3, "A3_1"), (4, "A4_1")).toDF("c1","c2") val r2 = Seq((3, "A3_2"), (4, "A4_2"), (5, "A5_2"), (6, "A6_2")).toDF("c1","c2") val mergedDF = r1.join(r2, Seq("c1"), "full_outer") .select( col("c1"), coalesce(r2("c2"), r1("c2")).alias("c2") // 优先用r2的c2,没有的话 fallback 到r1的c2 ) .orderBy("c1") mergedDF.show()
这个方法的优势在于扩展性强——如果后续新增字段,你可以给每个字段单独设置取值规则,比如某些字段保留r1的值,某些字段优先用r2的。
方法3:Spark 3.1+ 专属的replace方法(最简洁方案)
如果你使用的是Spark 3.1及以上版本,可以直接用DataFrame.replace方法,它专门用来用另一个DF替换当前DF中匹配键的行,再补充第二个DF里的新增行即可:
val r1 = Seq((1, "A1_1"), (2, "A2_1"), (3, "A3_1"), (4, "A4_1")).toDF("c1","c2") val r2 = Seq((3, "A3_2"), (4, "A4_2"), (5, "A5_2"), (6, "A6_2")).toDF("c1","c2") // 先用r2替换r1中重复键的行,再加上r2里r1没有的新行 val replacedR1 = r1.replace(r2, "c1") val newRowsInR2 = r2.filter(not(col("c1").isin(r1.select("c1").as[Int].collect():_*))) val mergedDF = replacedR1.union(newRowsInR2).orderBy("c1") mergedDF.show()
这个方法代码最简洁,但仅限Spark 3.1及以上版本使用。
内容的提问来源于stack exchange,提问作者Sent
相关产品推荐
相关产品推荐

