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

