You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.21 15:27:21