PySpark如何实现带30天时间约束和条件判断的窗口函数分组计算
PySpark 分组字段填充逻辑修正方案
原有代码错误点
- 窗口定义同时指定了
rangeBetween和rowsBetween,两个边界规则冲突,后定义的行范围会覆盖时间范围规则,导致30天有效期限制完全失效 - 使用
max()聚合只能取字段排序后的最大值,无法满足取最近一次访问对应的非Bad分组的需求 - 缺失30天内无有效非Bad值时回退为Bad的逻辑
- 窗口定义的语法不符合PySpark链式调用规范,缺少换行转义符或括号包裹
修正后完整代码
from pyspark.sql import functions as f from pyspark.sql.window import Window # 定义天数转秒的工具函数 days = lambda i: i * 86400 # 定义窗口:按用户分区,按日期转时间戳升序排序,范围为过去30天到当前行之前 user_30d_tracker = Window\ .partitionBy("userid")\ .orderBy(f.col("date").cast("timestamp").cast("long"))\ .rangeBetween(-days(30), -1) # 计算新分组字段 df = df.withColumn("Grouping(new)", f.when(f.col("Grouping(old)") != "Bad", f.col("Grouping(old)")) # 非Bad直接取原值 .otherwise( f.last( f.when(f.col("Grouping(old)") != "Bad", f.col("Grouping(old)")), ignorenulls=True # 忽略空值,取最近的有效分组 ).over(user_30d_tracker) ) .otherwise("Bad") # 30天内无有效分组则回退为Bad )
逻辑说明
- 窗口仅保留
rangeBetween时间范围规则,-days(30)对应30天前的时间戳,-1代表排除当前行,只取历史数据 - 使用
last()函数配合ignorenulls=True,可以拿到窗口范围内时间最近的非Bad分组值 - 两层
otherwise分别处理非Bad行直接赋值、无有效历史值回退为Bad的需求
内容的提问来源于stack exchange,提问作者Janice
相关产品推荐
相关产品推荐

