在PySpark中获取每个id的最大值
PySpark 实现按 ID 取最大值
你的需求是从包含id和数值列的数据集里,计算每个id对应的最大值,输出每个id及其对应的最大数值。下面提供两种常用实现方法:
方法一:使用 groupBy + agg 聚合函数
这是最直接的聚合方式,适合只需要分组后取最大值的场景。
示例代码
from pyspark.sql import SparkSession from pyspark.sql.functions import max # 初始化SparkSession spark = SparkSession.builder.appName("MaxByID").getOrCreate() # 模拟输入数据(匹配你的源数据格式) data = [("1", 10), ("1", 20), ("1", 15), ("2", 5), ("2", 8), ("3", 25)] df = spark.createDataFrame(data, ["id", "value"]) # 按id分组,计算每个id的最大值 max_df = df.groupBy("id").agg(max("value").alias("max_value")) # 展示结果 max_df.show()
输出结果
+---+----------+ | id|max_value| +---+----------+ | 1| 20| | 2| 8| | 3| 25| +---+----------+
方法二:使用窗口函数(适用于保留原表其他列的场景)
如果需要保留原数据中的其他字段,同时标记出每个id的最大值,窗口函数会更灵活。
示例代码
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import row_number spark = SparkSession.builder.appName("MaxByIDWindow").getOrCreate() data = [("1", 10), ("1", 20), ("1", 15), ("2", 5), ("2", 8), ("3", 25)] df = spark.createDataFrame(data, ["id", "value"]) # 定义窗口:按id分区,按value降序排序 window_spec = Window.partitionBy("id").orderBy(df["value"].desc()) # 添加行号列,取每个分区的第一行(即最大值所在行) ranked_df = df.withColumn("rank", row_number().over(window_spec)) max_df = ranked_df.filter(ranked_df["rank"] == 1).drop("rank") max_df.show()
输出结果
+---+-----+ | id|value| +---+-----+ | 1| 20| | 2| 8| | 3| 25| +---+-----+
内容的提问来源于stack exchange,提问作者sparc
相关产品推荐
相关产品推荐

