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

Pyspark如何实现基于用户最后一次行为倒推的30天tumbling window数据聚合

问题根因

PySpark 内置的F.window()函数默认是基于全局时间纪元对齐划分滚动窗口,不会针对每个用户的最后行为日期单独计算窗口边界,因此不符合你的个性化分窗需求。

实现方案

你可以按照「用户级最后日期打标→分窗分配→补全空窗口→聚合计算」的逻辑调整代码,具体实现如下:

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

# 步骤1:给每条数据拼接所属用户的最后行为日期
user_last_df = dfsp.withColumn("user_last_date", 
    F.max("date").over(W.partitionBy("id"))
)

# 步骤2:计算每条数据所属的30天倒推窗口边界
windowed_df = user_last_df.withColumn("days_diff", 
    F.datediff(F.col("user_last_date"), F.col("date"))
).withColumn("window_idx", 
    F.floor(F.col("days_diff") / 30)
).withColumn("window_end", 
    F.date_sub(F.col("user_last_date"), F.col("window_idx") * 30)
).withColumn("window_start", 
    F.date_sub(F.col("window_end"), 30)
)

# 步骤3:生成每个用户的完整窗口序列,补全无数据的空窗口
user_window_range = windowed_df.groupBy("id", "user_last_date").agg(
    F.max("window_idx").alias("max_window_idx")
)
# 如果需要固定给所有用户生成最近N个30天窗口(比如全年12个),可以把下方sequence的第二个参数改成F.lit(11)
user_all_windows = user_window_range.withColumn("window_idx", 
    F.explode(F.sequence(F.lit(0), F.col("max_window_idx")))
).withColumn("window_end", 
    F.date_sub(F.col("user_last_date"), F.col("window_idx") * 30)
).withColumn("window_start", 
    F.date_sub(F.col("window_end"), 30)
).select("id", "window_start", "window_end")

# 步骤4:关联原始数据聚合求和,空窗口填充0
result_df = user_all_windows.join(
    windowed_df,
    on=["id", "window_start", "window_end"],
    how="left"
).groupBy("id", "window_start", "window_end").agg(
    F.coalesce(F.sum("value"), F.lit(0)).alias("sum(value)")
).withColumn("window", 
    F.concat(F.lit("("), F.col("window_start"), F.lit(", "), F.col("window_end"), F.lit(")"))
).select("id", "window", "sum(value)").orderBy("id", "window_start")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 17:15:07