如何用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
相关产品推荐
相关产品推荐

