Scala Spark如何提取DataFrame数组列最小值并新增列

问题说明
现有上述示例DataFrame,需求为提取key列中的最小数值,新增一列展示该最小值。原有Scala实现代码无法正常运行,以下是错误说明和修正方案。
错误原因
- UDF参数类型错误:原代码写的
dat1[String]是DataFrame变量名,并非合法的参数类型。Spark中存储字符串数组的列,在UDF中需要用Seq[String]或Array[String]接收。 - 语法括号不匹配:原代码漏写了
select方法的闭合右括号,错误将show方法写在了as方法的参数列表内部,不符合Scala语法规则。
可直接运行的修正代码
import org.apache.spark.sql.functions._ // 定义UDF:过滤空字符串、转整数后取最小值,空数组返回0 val minUdf = udf((arr: Seq[String]) => { val filtered = arr.filterNot(_ == "") if(filtered.isEmpty) 0 else filtered.map(_.toInt).min }) // 修正括号匹配,查询后展示结果 dat1.select(col("placekey"), minUdf(col("key")).as("lowest_value")).show(false)
性能优化方案(无UDF实现)
大数据量场景不推荐使用自定义UDF,UDF无法利用Spark原生Catalyst优化器,性能远低于内置函数。可直接用Spark内置数组函数实现相同逻辑:
import org.apache.spark.sql.functions._ dat1.select( col("placekey"), coalesce( array_min(transform(filter(col("key"), _ =!= ""), _.cast("int"))), lit(0) ).as("lowest_value") ).show(false)
内容的提问来源于stack exchange,提问作者Tammy Nolasco
相关产品推荐
相关产品推荐

