Scala DataFrame去重:优先保留DOB非空行,无则留空行
解决方案
针对你的需求,不需要复杂拆分DataFrame再Union,有两种高效实现方式:
方法一:窗口函数(推荐,代码简洁且性能更优)
利用窗口函数给每个AccNo分组内的行排序,让非空DOB的行优先级更高,然后取每组第一行即可:
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ // 定义窗口:按AccNo分区,非空DOB的行排前面 val windowSpec = Window.partitionBy("AccNo") .orderBy(when($"DOB".isNotNull, 1).otherwise(2)) // 添加行号,过滤取每组第一行后删除行号列 val resultDF = df.withColumn("row_num", row_number().over(windowSpec)) .filter($"row_num" === 1) .drop("row_num")
方法二:优化后的Union方式
如果坚持用Union思路,可以避免冗余计算,步骤如下:
- 先提取所有存在非空
DOB的AccNo - 分别获取「非空DOB的行」和「无对应非空DOB的空行」,再合并
// 获取所有有非空DOB的AccNo val validAccNos = df.filter($"DOB".isNotNull).select("AccNo").distinct() // 第一部分:保留所有非空DOB的行 val nonNullDF = df.filter($"DOB".isNotNull) // 第二部分:保留那些没有非空DOB记录的AccNo对应的空行 val nullOnlyDF = df.filter($"DOB".isNull) .join(validAccNos, Seq("AccNo"), "left_anti") // 合并两个结果 val resultDF = nonNullDF.union(nullOnlyDF)
两种方法都能得到你想要的输出,其中窗口函数只需扫描一次数据,大表场景下更推荐使用。
内容的提问来源于stack exchange,提问作者ShaanGreat
相关产品推荐
相关产品推荐

