如何为DataFrame重复值检测函数实现ignoreNulls参数?
解决方案
修改后的函数代码
import org.apache.spark.sql.functions._ import org.apache.spark.sql.Window def checkRepeatedKey(newColName: String, keys: Seq[String], ignoreNulls: Boolean = false)(dataframe: DataFrame): DataFrame = { // 判断当前行的指定keys列是否包含null值 val hasNullInKeys = keys.map(col(_).isNull).reduce(_ || _) // 按指定keys列分组的窗口条件 val windowCondition = Window.partitionBy(keys: _*) // 计算每组的行数 val sumCount = sum(lit(1)).over(windowCondition) dataframe .withColumn("sum_count", sumCount) .withColumn(newColName, // 分支处理:ignoreNulls为true且行含null时直接返回false,否则判断是否重复 when(ignoreNulls && hasNullInKeys, lit(false)) .otherwise(col("sum_count") > 1) ) .drop("sum_count") }
关键逻辑说明
- 参数新增:添加默认值为
false的ignoreNulls参数,匹配需求中的默认行为。 - null值判断:通过
hasNullInKeys表达式一键判断当前行的指定keys列是否存在null,避免逐行硬编码判断。 - 窗口统计:窗口函数
partitionBy(keys: _*)会自动将null值归为一组(Spark原生行为),满足ignoreNulls=false时的重复判定需求。 - 分支处理:用Spark内置的
when函数实现行级条件判断:- 当
ignoreNulls=true且当前行含null时,直接标记为false(忽略null的重复统计) - 其他场景(
ignoreNulls=false或行无null),按分组行数是否大于1判定重复
- 当
测试验证
测试默认场景(ignoreNulls=false)
checkRepeatedKey("is_repeated", Seq("name"))(testDF).show()
输出结果:
+--------+--------------+-----------+ |name_key| name|is_repeated| +--------+--------------+-----------+ | 1| name-1| false| | 2|repeated-name| true| | 3|repeated-name| true| | 4| name-4| false| | 5| null| true| | 6| null| true| +--------+--------------+-----------+
测试忽略null场景(ignoreNulls=true)
checkRepeatedKey("is_repeated", Seq("name"), ignoreNulls = true)(testDF).show()
输出结果:
+--------+--------------+-----------+ |name_key| name|is_repeated| +--------+--------------+-----------+ | 1| name-1| false| | 2|repeated-name| true| | 3|repeated-name| true| | 4| name-4| false| | 5| null| false| | 6| null| false| +--------+--------------+-----------+
内容的提问来源于stack exchange,提问作者Malkath
相关产品推荐
相关产品推荐

