PySpark:分组获取上月最后值,避免自连接实现方案
仅用窗口函数实现分组获取上月最后属性值
需求是按id分组,为每条记录添加上个月对应分组的最后一条value值,原实现采用自连接方式,效率较低,下面提供纯窗口函数的优化方案。
实现代码
from pyspark.sql import SparkSession, functions as F from pyspark.sql.window import Window spark = SparkSession.builder.appName("prev_month_last_value").getOrCreate() df = ( spark.createDataFrame( [ ("2022-07-29", 1, 1), ("2022-07-30", 1, 2), ("2022-07-31", 1, 3), ("2022-08-01", 1, 4), ("2022-08-02", 1, 5), ("2022-08-03", 1, 6), ("2022-09-10", 1, 8), ("2022-09-11", 1, 9), ("2022-09-12", 1, 10), ("2022-07-29", 2, 7), ("2022-07-30", 2, 6), ("2022-07-31", 2, 5), ("2022-08-01", 2, 4), ("2022-08-02", 2, 3), ("2022-08-03", 2, 2), ("2022-09-10", 2, 8), ("2022-09-11", 2, 9), ("2022-09-12", 2, 10), ], ["date","id","value"] ) .withColumn("date", F.to_date(F.col("date"))) ) # 1. 计算每条记录的当前月份和上个月标识 df_with_months = df.withColumn( "current_month", F.date_trunc("month", F.col("date")) ).withColumn( "prev_month", F.add_months(F.col("current_month"), -1) ) # 2. 定义窗口:按id分组,按日期升序,覆盖分组内从第一行到当前行的所有数据 w = Window.partitionBy("id").orderBy("date").rowsBetween(Window.unboundedPreceding, Window.currentRow) # 3. 收集历史数据并筛选上月最后一条value df_result = df_with_months.withColumn( # 收集到当前行为止的所有(月份, 值)结构,后续反转取最新数据 "month_value_list", F.collect_list(F.struct("current_month", "value")).over(w) ).withColumn( # 反转列表后筛选出上月的条目,取第一个即为上月最后一条value "last_value_prev_month", F.array_first( F.filter( F.reverse(F.col("month_value_list")), lambda x: x.current_month == F.col("prev_month") ) ).value ).drop("current_month", "prev_month", "month_value_list").orderBy("id", "date") df_result.show()
输出结果
+---+----------+-----+---------------------+ | id| date|value|last_value_prev_month| +---+----------+-----+---------------------+ | 1|2022-07-29| 1| null| | 1|2022-07-30| 2| null| | 1|2022-07-31| 3| null| | 1|2022-08-01| 4| 3| | 1|2022-08-02| 5| 3| | 1|2022-08-03| 6| 3| | 1|2022-09-10| 8| 6| | 1|2022-09-11| 9| 6| | 1|2022-09-12| 10| 6| | 2|2022-07-29| 7| null| | 2|2022-07-30| 6| null| | 2|2022-07-31| 5| null| | 2|2022-08-01| 4| 5| | 2|2022-08-02| 3| 5| | 2|2022-08-03| 2| 5| | 2|2022-09-10| 8| 2| | 2|2022-09-11| 9| 2| | 2|2022-09-12| 10| 2| +---+----------+-----+---------------------+
核心逻辑说明
- 先通过
date_trunc和add_months计算每条记录的当前月份与上个月,为后续匹配做准备。 - 窗口函数按
id分组、日期升序排列,确保能收集到当前行之前的所有历史数据。 - 利用
collect_list聚合历史的月份和值,反转列表后筛选上月数据,取第一个元素就是上月最后一条记录的value(反转后最新数据排在最前)。 - 全程无需自连接,减少了数据shuffle的开销,执行效率更高。
内容的提问来源于stack exchange,提问作者Sebastian Hohmann
相关产品推荐
相关产品推荐

