You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.28 10:54:03