PySpark DataFrame按日提取各ID最新更新记录的实现方法
解决方案
要实现每个ID每日仅保留最后更新记录的需求,核心是基于ID和stock_date分组,在每个分组内筛选出updated_at最晚的行。可以通过PySpark的窗口函数完成,具体步骤如下:
1. 导入必要函数
from pyspark.sql import Window from pyspark.sql.functions import row_number, desc
2. 定义窗口规则
创建窗口,按ID和stock_date分区(确保同一ID同一日期的记录归为一组),并按updated_at降序排序(让最新更新的记录排在每组首位):
window_spec = Window.partitionBy("ID", "stock_date").orderBy(desc("updated_at"))
3. 添加分组内的行号
为每条记录添加分组内行号,最新更新的记录行号为1:
df_with_row_num = df.withColumn("daily_rank", row_number().over(window_spec))
4. 筛选目标记录
仅保留每个分组内行号为1的记录,即每个ID每日的最新更新行:
result_df = df_with_row_num.filter(df_with_row_num.daily_rank == 1)
结果验证
执行上述代码后,输出结果将匹配你预期的行(row_num为1、3、5、6、7、9、10、11、13、15),每个ID的每个日期仅保留最后更新的记录,不同日期的记录全部保留。
内容的提问来源于stack exchange,提问作者Leandro Destefani
相关产品推荐
相关产品推荐

