PySpark中如何在groupBy后按最大日期过滤数据
解决方法
你的代码无法运行的原因是:groupBy操作会将DataFrame转换为聚合后的结果集,此时原DataFrame中的date列不再保留每行的原始值,无法直接与聚合函数max(date)进行比较判断。
要实现每个公司每年保留日期最新的一行数据,最直接高效的方式是使用窗口函数,具体步骤如下:
步骤1:导入必要模块
from pyspark.sql import functions as f from pyspark.sql.window import Window
步骤2:定义窗口规则
按company_id和date_year分区(即分组),在每个分区内按date字段降序排序:
window_spec = Window.partitionBy("company_id", "date_year").orderBy(f.desc("date"))
步骤3:添加行号并过滤目标行
给每行添加分区内的行号,日期最大的行在分区内的行号为1,最后过滤出行号等于1的记录即可:
df_orders = df_orders \ .withColumn("row_num", f.row_number().over(window_spec)) \ .filter(f.col("row_num") == 1) \ .drop("row_num") # 移除临时生成的行号列
补充说明
- 如果同一个分组内存在多条日期相同的最大日期记录,
row_number()会随机保留其中一条;若想保留所有日期最大的记录,可以改用rank()或dense_rank()替换row_number()。 - 另一种替代方案是先聚合得到每个分组的最大日期,再与原DataFrame关联,但窗口函数在大数据量场景下的性能通常更优。
内容的提问来源于stack exchange,提问作者Andrii
相关产品推荐
相关产品推荐

