使用Spark DataFrame找出列中排除0的最小值并新增列
Spark DataFrame 计算每行排除0.0后的最小值
你用least()函数得不到预期结果,是因为least()会把0.0纳入计算范围,比如第二行least(200.0, 0.0, 0.0)会返回0.0,不符合你排除0.0的需求。
下面是两种高效的解决方案,用Spark内置函数实现,无需自定义UDF:
方法1:使用SQL表达式(简洁直观)
通过array()将每行的列转为数组,用filter()剔除0.0元素,再用array_min()取剩余元素的最小值:
Python 代码
from pyspark.sql import SparkSession from pyspark.sql.functions import expr # 创建示例DataFrame spark = SparkSession.builder.appName("ExcludeZeroMin").getOrCreate() data = [(100.0, 120.0, 150.0), (200.0, 0.0, 0.0), (0.0, 20.0, 100.0)] df = spark.createDataFrame(data, ["Column1", "Column2", "Column3"]) # 新增Least列 df_with_least = df.withColumn( "Least", expr("array_min(filter(array(Column1, Column2, Column3), x -> x != 0.0))") ) # 查看结果 df_with_least.show()
Scala 代码
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions.expr // 创建示例DataFrame val spark = SparkSession.builder.appName("ExcludeZeroMin").getOrCreate() val data = Seq((100.0, 120.0, 150.0), (200.0, 0.0, 0.0), (0.0, 20.0, 100.0)) val df = spark.createDataFrame(data).toDF("Column1", "Column2", "Column3") // 新增Least列 val dfWithLeast = df.withColumn( "Least", expr("array_min(filter(array(Column1, Column2, Column3), x -> x != 0.0))") ) // 查看结果 dfWithLeast.show()
方法2:使用Spark函数组合
用when()标记非0.0的值,转为数组后过滤掉null,再取最小值:
Python 代码
from pyspark.sql.functions import array, array_min, col, when df_with_least = df.withColumn( "Least", array_min( array( when(col("Column1") != 0.0, col("Column1")), when(col("Column2") != 0.0, col("Column2")), when(col("Column3") != 0.0, col("Column3")) ).filter(lambda x: x.isNotNull()) ) ) df_with_least.show()
处理全0.0的特殊情况
如果某行所有列都是0.0,array_min()会返回null。如果需要返回0.0,可以用coalesce()兜底:
from pyspark.sql.functions import coalesce, lit df_with_least = df.withColumn( "Least", coalesce( expr("array_min(filter(array(Column1, Column2, Column3), x -> x != 0.0))"), lit(0.0) ) )
预期输出
运行上述代码后,得到的DataFrame如下:
| Column1 | Column2 | Column3 | Least |
|---|---|---|---|
| 100.0 | 120.0 | 150.0 | 100.0 |
| 200.0 | 0.0 | 0.0 | 200.0 |
| 0.0 | 20.0 | 100.0 | 20.0 |
内容的提问来源于stack exchange,提问作者RMK
相关产品推荐
相关产品推荐

