如何在PySpark DataFrame中按行取多列的最小值?
在PySpark DataFrame中计算指定列的逐行最小值
没问题,我来帮你搞定这个需求!先明确一下你的场景:
原始DataFrame结构如下:
- - - - - - A| B | C| - - - - - - 1| 2 | 3| 2| 3 | 6| 3| 5 | 4|
你需要按行计算B、C两列的最小值,最终得到只保留A列和最小值列的结果:
- - - - - - A| min(B,C) - - - - - - 1| 2 2| 3 3| 4
具体实现步骤
PySpark内置的least()函数正好能解决这个问题——它可以逐行比较多个列的值,返回每行中的最小值。下面是完整的代码示例:
- 导入必要模块
from pyspark.sql import SparkSession from pyspark.sql.functions import least
- 初始化SparkSession并创建示例DataFrame
# 启动SparkSession实例 spark = SparkSession.builder.appName("RowwiseMinCalculation").getOrCreate() # 构造你提供的原始数据 raw_data = [(1, 2, 3), (2, 3, 6), (3, 5, 4)] original_df = spark.createDataFrame(raw_data, schema=["A", "B", "C"]) # 可选:查看原始数据确认结构 original_df.show()
- 计算逐行最小值并生成目标结果
# 添加新列,计算B和C的最小值 temp_df = original_df.withColumn("min(B,C)", least(original_df["B"], original_df["C"])) # 只保留需要的A列和新生成的最小值列 result_df = temp_df.select("A", "min(B,C)") # 查看最终结果 result_df.show()
补充小提示
- 如果需要计算超过两列的逐行最小值,直接给
least()函数传入更多列即可,比如least(col("B"), col("C"), col("D"))。 - 如果你的列名包含特殊字符或空格,记得用
col()函数包裹列名,比如least(col("B value"), col("C value"))。
内容的提问来源于stack exchange,提问作者wrek
相关产品推荐
相关产品推荐

