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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 04:01:07