Spark Scala实现无硬编码合并两单条记录DataFrame(保留非空值)
Spark Scala 合并单条记录DataFrame并保留非空值(自动适配所有列)
需求说明
现有两个各包含单条记录的DataFrame(DF1、DF2),需合并二者并保留每列的非空值;由于涉及200+列,禁止硬编码列名。
示例数据
DF1:
+---+-----+-----+-----+-----+ | id| sid1| pid1| sid2| pid2| +---+-----+-----+-----+-----+ | 1| 1111| null| 2222| null| +---+-----+-----+-----+-----+
DF2:
+---+-----+-----+-----+-----+ | id| sid1| pid1| sid2| pid2| +---+-----+-----+-----+-----+ | 1| null| 3333| null| 4444| +---+-----+-----+-----+-----+
期望输出
+---+-----+-----+-----+-----+ | id| sid1| pid1| sid2| pid2| +---+-----+-----+-----+-----+ | 1| 1111| 3333| 2222| 4444| +---+-----+-----+-----+-----+
解决方案
利用Spark的coalesce函数(返回第一个非空值),结合动态列名遍历实现无硬编码的合并:
步骤1:导入依赖
import org.apache.spark.sql.functions._ import org.apache.spark.sql.Column
步骤2:获取所有列名
从任意一个DataFrame中提取列名(确保两DF列名一致,若不一致可取列名交集):
val allColumns = df1.columns
步骤3:动态生成合并列表达式
对每一列生成coalesce表达式,优先取DF1的非空值,DF1为空则取DF2的值:
val coalesceColumns: Array[Column] = allColumns.map(colName => coalesce(df1(colName), df2(colName)).alias(colName) )
步骤4:合并DataFrame
由于两DF各含单条记录,可直接用crossJoin合并,再选择动态生成的列:
val mergedDF = df1.crossJoin(df2).select(coalesceColumns: _*)
如果存在主键(如示例中的id),用主键关联更严谨:
val mergedDF = df1.join(df2, Seq("id"), "inner").select(coalesceColumns: _*)
验证结果
执行show()查看合并后的数据:
mergedDF.show()
原理说明
coalesce函数会依次检查传入的列值,返回第一个非空值,完美匹配"保留非空值"的需求;- 通过遍历列名动态生成表达式,无需硬编码任何列名,可适配任意数量的列(包括200+列的场景)。
内容的提问来源于stack exchange,提问作者Ankit Gupta
相关产品推荐
相关产品推荐

