PySpark中按类别获取对应最大值的日期时间并添加至原数据集
PySpark中按类别获取对应最大值的日期时间并添加至原数据集
嗨,我来帮你搞定这个需求!你想要给数据集中的每一行都添加一列datetimeMax,用来显示该行所属类别下value最大值对应的日期时间,对吧?下面我用PySpark给你两种实用的实现方案:
先构建示例数据集
首先咱们先把你给出的示例数据转换成PySpark的DataFrame,方便后续操作:
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("MaxDatetimePerCategory").getOrCreate() # 构建你提供的示例数据 data = [ ("a", "date1", 10), ("a", "date2", 30), ("a", "date3", 20), ("a", "date4", 50), ("a", "date5", 30), ("b", "date6", 20), ("b", "date7", 15), ("b", "date8", 30), ("b", "date9", 40), ("c", "date10", 10), ("c", "date11", 10), ("c", "date12", 30) ] df = spark.createDataFrame(data, ["category", "datetime", "value"]) df.show()
方案一:分组聚合+关联(适合新手理解)
这个方案分三步走,逻辑清晰,容易上手:
- 先按
category分组,计算每个类别下的最大value - 关联原数据集,找到每个类别对应最大
value的日期时间(如果有多个相同最大值,这里取第一个出现的日期) - 把得到的
datetimeMax关联回原数据集,得到最终结果
# 步骤1:获取每个category的最大value max_value_df = df.groupBy("category").agg(F.max("value").alias("max_value")) # 步骤2:找到每个category对应max_value的datetime category_max_datetime = df.join(max_value_df, on=["category", "value"]) \ .groupBy("category") \ .agg(F.first("datetime").alias("datetimeMax")) # 步骤3:关联回原数据集,添加datetimeMax列 result_df = df.join(category_max_datetime, on="category") result_df.show()
方案二:窗口函数(更简洁高效)
用窗口函数可以一步到位,不需要多次关联,代码更简洁:
我们定义一个按category分区的窗口,先按value降序排序(这样每个组里第一行就是value最大的行),然后用first()函数提取该行的datetime作为整个组的datetimeMax。
# 定义窗口:按category分区,按value降序排列(若有相同最大值,可追加datetime排序控制取哪一个) window_spec = Window.partitionBy("category").orderBy(F.desc("value")).rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing) # 直接添加datetimeMax列 result_df = df.withColumn("datetimeMax", F.first("datetime").over(window_spec)) result_df.show()
如果你的数据中存在同一个类别下多个行的value都是最大值的情况,想要取最后出现的日期,只需要把排序规则改成F.desc("value"), F.desc("datetime")即可。
运行上面的代码后,就能得到你想要的结果啦!
备注:内容来源于stack exchange,提问作者Nabs335
相关产品推荐
相关产品推荐

