如何使用Pyspark按指定字段分组取最大日期并返回所有符合条件记录
Pyspark 实现分组取每组最大日期对应全量记录
实现思路
使用窗口函数对userId和memberId分组,分组内按date降序排序并标记排名,筛选出排名为1的所有记录即为每组最大日期对应的全部数据。使用rank()函数而非row_number(),可保证同一分组下多条最大日期记录都能被筛选出来。
完整实现代码
from pyspark.sql import SparkSession from pyspark.sql.window import Window import pyspark.sql.functions as F # 初始化SparkSession spark = SparkSession.builder.appName("fetch_max_date_records").getOrCreate() # 构建样例数据集 source_data = [ ("2016-04-06", 1234, 111, 1), ("2016-04-06", 1234, 222, 5), ("2016-04-06", 1234, 111, 8), ("2016-04-06", 1234, 222, 9), ("2016-04-05", 4567, 111, 1), ("2016-04-06", 4567, 222, 5), ("2016-04-06", 4567, 111, 8), ("2016-04-06", 4567, 222, 9) ] df = spark.createDataFrame(source_data, schema=["date", "userId", "memberId", "value"]) # 定义窗口:按userId、memberId分组,分组内按日期倒序排序 window_rule = Window.partitionBy("userId", "memberId").orderBy(F.desc("date")) # 标记每条记录在分组内的排名 df_with_rank = df.withColumn("rank", F.rank().over(window_rule)) # 筛选排名为1的记录,重命名date字段为datetime result_df = df_with_rank.filter(F.col("rank") == 1)\ .select("userId", "memberId", F.col("date").alias("datetime"), "value") # 打印结果 result_df.show()
输出结果调整说明
上述基础实现会输出每个分组下所有最大日期的原始记录,如果需要匹配你给出的预期输出(同一最大日期下取最大value),可以补充分组聚合逻辑:
result_df = df_with_rank.filter(F.col("rank") == 1)\ .groupBy("userId", "memberId", F.col("date").alias("datetime"))\ .agg(F.max("value").alias("value")) result_df.show()
聚合后输出完全匹配预期结果:
+------+--------+----------+-----+ |userId|memberId| datetime|value| +------+--------+----------+-----+ | 1234| 111|2016-04-06| 8| | 1234| 222|2016-04-06| 9| | 4567| 111|2016-04-06| 8| | 4567| 222|2016-04-06| 9| +------+--------+----------+-----+
内容的提问来源于stack exchange,提问作者Pavithra Kannan
相关产品推荐
相关产品推荐

