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

如何用PySpark Window实现按id分组的30天滚动求和及去重日期统计?

PySpark实现30天滚动求和与去重日期计数

问题背景

现有包含date、id、value三列的PySpark DataFrame,示例数据如下:

df_tmp = spark.createDataFrame([('2023-01-01', 1001, 5),
                                ('2023-01-15', 1001, 3),
                                ('2023-02-10', 1001, 1),
                                ('2023-02-20', 1001, 2),
                                ('2023-01-02', 1002, 7),
                                ('2023-01-02', 1002, 6),
                                ('2023-01-03', 1002, 1)],
                               ["date", "id", "value"])

需要按id分组,对每条记录计算两个指标:

  • 30天滚动求和:当前日期过去30天内(不含当前记录的value)的value总和
  • 过去30天内去重日期数:当前日期过去30天内(不含当前记录的日期)该id出现的不同日期数量

期望输出:

+----------+----+-----+----------------+-------------------------+
|      date|  id|value|30_day_value_sum|days_seen_in_past_30_days|
+----------+----+-----+----------------+-------------------------+
|2023-01-01|1001|    5|               0|                        0|
|2023-01-15|1001|    3|               0|                        1|
|2023-02-10|1001|    1|               3|                        1|
|2023-02-20|1001|    2|               1|                        2|
|2023-01-02|1002|    7|               0|                        0|
|2023-01-02|1002|    6|               7|                        1|
|2023-01-03|1002|    1|              13|                        2|
+----------+----+-----+----------------+-------------------------+

解决方案

通过PySpark窗口函数实现,步骤如下:

1. 导入依赖并转换日期类型

先将字符串类型的date列转为日期类型,方便后续日期范围计算:

from pyspark.sql import functions as F
from pyspark.sql.window import Window

# 转换date列为日期类型
df = df_tmp.withColumn("date", F.to_date("date"))

2. 定义滚动窗口

按id分区,按date排序,窗口范围设置为当前日期前30天到当前行的前一行(确保不包含当前记录的数值和日期):

# 将日期转为时间戳(秒),用于rangeBetween的数值范围计算
window_spec = Window.partitionBy("id") \
    .orderBy(F.col("date").cast("long")) \
    .rangeBetween(-30 * 86400, -1)  # 86400为一天的秒数,-30*86400表示往前推30天

3. 计算目标指标

使用窗口函数计算滚动求和与去重日期数,并用coalesce处理窗口无数据时的null值,转为0:

df_result = df.withColumn("30_day_value_sum", F.coalesce(F.sum("value").over(window_spec), F.lit(0))) \
              .withColumn("days_seen_in_past_30_days", F.coalesce(F.countDistinct("date").over(window_spec), F.lit(0)))

4. 查看结果

df_result.show()

执行后即可得到符合要求的输出。

内容的提问来源于stack exchange,提问作者Abhishek Parab

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 01:09:55