Spark UDF如何传入列数组统计DataFrame空值列数量?
问题:统计DataFrame中值为"NULL"的列数量(UDF传参问题)
我有一个包含N列的DataFrame,想要新增一列以统计其中值为"NULL"的列的数量。我尝试创建UDF实现该功能,但因无法传入列数组参数导致代码无法正常运行,示例代码如下:
val simpleData = Seq( ("row1", "NULL" , "NULL" , "NULL" , "NULL" , "NULL", "1"), ("row2", "1", "NULL", "2023", "NULL", "01", "NULL")) val myDs = simpleData.toDF("row", "field1", "field2", "field3", "field4", "field5", "field6") myDs.show() val windowcols = myDs.columns.filterNot(List("row").contains(_)) def countNullsUDF: UserDefinedFunction = udf { (values: List[String]) => values.filter( value => value == "NULL").length } myDs.withColumn("columnsWithNull", countNullsUDF(windowcols)).show(10, false)
请问是否可以向该UDF传入列数组或类似参数?
解决方案
方法1:修改UDF传参方式,结合array()函数
不能直接把列名数组传给UDF,因为UDF需要接收的是每行中这些列的值集合而非列名。需要用Spark的array()函数将目标列打包成数组类型的Column,再传给UDF:
import org.apache.spark.sql.functions.{col, udf, array} import org.apache.spark.sql.UserDefinedFunction val simpleData = Seq( ("row1", "NULL" , "NULL" , "NULL" , "NULL" , "NULL", "1"), ("row2", "1", "NULL", "2023", "NULL", "01", "NULL")) val myDs = simpleData.toDF("row", "field1", "field2", "field3", "field4", "field5", "field6") val windowcols = myDs.columns.filterNot(List("row").contains(_)) // 调整UDF参数类型为Seq[String],兼容Spark array()返回的WrappedArray def countNullsUDF: UserDefinedFunction = udf { (values: Seq[String]) => values.count(_ == "NULL") } // 用array()将列数组转为数组列后传入UDF myDs.withColumn("columnsWithNull", countNullsUDF(array(windowcols.map(col): _*))) .show(10, false)
运行后会得到正确结果:
+----+------+------+------+------+------+------+---------------+ |row |field1|field2|field3|field4|field5|field6|columnsWithNull| +----+------+------+------+------+------+------+---------------+ |row1|NULL |NULL |NULL |NULL |NULL |1 |5 | |row2|1 |NULL |2023 |NULL |01 |NULL |3 | +----+------+------+------+------+------+------+---------------+
方法2:使用Spark内置函数(更高效,推荐)
无需自定义UDF,直接用sum+when组合实现,避免UDF的序列化开销,性能更优:
import org.apache.spark.sql.functions.{col, sum, when} // 对每个目标列判断是否为"NULL",是则记1否则记0,最后求和 val nullCountCol = windowcols.map(c => when(col(c) === "NULL", 1).otherwise(0)) .reduce(_ + _) .alias("columnsWithNull") myDs.withColumn("columnsWithNull", nullCountCol) .show(10, false)
这段代码实现的效果和UDF完全一致,但执行效率更高。
内容的提问来源于stack exchange,提问作者Guille
相关产品推荐
相关产品推荐

