Scala通用函数:用DataFrame2值填充DataFrame1空值(去重列)
Spark多列空值替换通用实现方案
需求说明
现有两个列名完全一致的DataFrame(各含100列),需实现:
- 将DataFrame1(左表)所有列的空值替换为DataFrame2(右表)对应列的值
- 若DataFrame2对应值也为空,则保留DataFrame1的原值
现有单列方案的问题
当前仅实现了单列处理逻辑,存在两个明显缺陷:
- 无法批量处理100列,扩展性差
- 执行后会生成重复列,需要额外清理
单列实现代码及结果如下:
val df5 = Seq((1, null), (2, "B")).toDF("id", "colA") val df6 = Seq((1, "C"), (4, "D")).toDF("id", "colA")
+---+----+ | id|colA| +---+----+ --> df5.show | 1|null| | 2| B| +---+----+ +---+----+ | id|colA| +---+----+ | 1| C| --> df6.show | 4| D| +---+----+
import org.apache.spark.sql._ import org.apache.spark.sql.functions._ def replaceNull(A:DataFrame, B:DataFrame) :DataFrame = { A.join(B, Seq("id"), "left") .withColumn("colA", when(A("colA").isNull, B("colA")).otherwise(A("colA"))) } replaceNull(df5, df6).show
+---+----+----+ | id|colA|colA| +---+----+----+ | 1| C| C| | 2| B| B| +---+----+----+
通用解决方案
以下是支持所有列批量处理的Scala函数,自动完成空值替换并避免重复列:
import org.apache.spark.sql._ import org.apache.spark.sql.functions._ def replaceAllNulls(leftDF: DataFrame, rightDF: DataFrame, joinKey: String): DataFrame = { // 提取除关联键外的所有业务列名 val businessColumns = leftDF.columns.filter(_ != joinKey) // 执行左关联,保留左表全量数据 val joinedDF = leftDF.join(rightDF, Seq(joinKey), "left") // 对每个业务列生成空值替换逻辑,用原列名别名避免重复 val transformedColumns = businessColumns.map { colName => when(leftDF(colName).isNull, rightDF(colName)).otherwise(leftDF(colName)).alias(colName) } // 仅保留关联键和处理后的业务列 joinedDF.select(col(joinKey) +: transformedColumns: _*) }
使用示例
// 构造含多列的测试数据 val df5 = Seq((1, null, "X"), (2, "B", null)).toDF("id", "colA", "colB") val df6 = Seq((1, "C", "Y"), (4, "D", "Z")).toDF("id", "colA", "colB") // 调用通用函数 val resultDF = replaceAllNulls(df5, df6, "id") resultDF.show()
输出结果
+---+----+----+ | id|colA|colB| +---+----+----+ | 1| C| X| | 2| B|null| +---+----+----+
代码逻辑说明
- 列过滤:自动排除关联键,仅处理业务列,适配任意数量的列
- 左关联:保证左表所有数据不丢失,右表仅补充匹配数据
- 批量处理:对每个列生成统一的空值替换规则,用原列名别名覆盖重复列
- 结果裁剪:只保留需要的列,自动剔除关联后产生的重复列
内容的提问来源于stack exchange,提问作者SparkLearner
相关产品推荐
相关产品推荐

