Spark Scala中DataFrame多列聚合:获取指定列非空最值
Scala Spark 分组聚合动态指定列取非空最值
问题场景
现有如下Spark DataFrame(实际包含30+列):
val df = Seq( ("GUID1", Some(1), Some(22), Some(30), Some(56)), ("GUID1", Some(4), None, Some(35), Some(52)), ("GUID1", None, Some(24), None, Some(58)), ("GUID2", Some(5), Some(21), Some(31), None) ).toDF("GUID", "A", "B", "C", "D" )
需求:按GUID分组后,对指定列计算非空最值:
- 数组
max_cols中的A、B列取非空最大值 - 数组
min_cols中的C、D列取非空最小值
解决方案
Spark的max和min函数默认会忽略null值,因此无需复杂的数组操作,直接动态生成聚合表达式即可:
import org.apache.spark.sql.functions.{max, min, col} val max_cols = Array("A", "B") val min_cols = Array("C", "D") // 生成max列的聚合表达式 val maxAggs = max_cols.map(c => max(col(c)).alias(c)) // 生成min列的聚合表达式 val minAggs = min_cols.map(c => min(col(c)).alias(c)) // 合并表达式并执行分组聚合 val resultDF = df.groupBy("GUID").agg(maxAggs ++ minAggs: _*) resultDF.show()
输出结果
+-----+---+---+---+----+ | GUID| A| B| C| D| +-----+---+---+---+----+ |GUID1| 4| 24| 30| 52| |GUID2| 5| 21| 31|null| +-----+---+---+---+----+
说明
- 动态生成表达式:通过遍历列名数组,批量生成
max/min聚合表达式,适配30+列的场景,避免重复手写代码 - 自动忽略null:Spark内置的
max和min函数在计算时会自动跳过null值,完全符合需求 - 类型安全:使用
col(c)确保列名引用的正确性,避免字符串拼接可能出现的错误
内容的提问来源于stack exchange,提问作者Ganesha
相关产品推荐
相关产品推荐

