PySpark DataFrame中筛选前续为连续5个0的首个1
嘿,这个需求我之前处理过,单纯用rank函数确实搞不定,因为它没法捕捉连续序列的状态变化。咱们可以用PySpark的窗口函数结合状态标记的方式来精准定位目标行,一步步来:
假设前提
先确认你的DataFrame包含这三列:
user_id: 用户唯一标识month_col: 月份列(必须是可排序的类型,比如date/timestamp或格式如'2023-01'的字符串)event: 取值0或1的目标列
解决方案步骤
1. 给用户内的行按时间排序并加行号
这一步是为了后续能精准定位前后行的位置:
from pyspark.sql import Window import pyspark.sql.functions as F # 定义用户内按时间排序的窗口 user_time_window = Window.partitionBy("user_id").orderBy("month_col") # 添加行号,方便后续判断位置 df = df.withColumn("row_num", F.row_number().over(user_time_window))
2. 直接检查目标模式:当前行是1,且恰好前面连续5个0
这里用窗口函数的collect_list抓取前面5行的event值,同时排除前面有更多0的情况(比如连续6个0后出现的1):
df = df.withColumn( "is_target", F.when( F.col("event") == 1, F.and( # 前面5行必须全是0 F.expr("array_join(collect_list(event) over (user_time_window rows between 5 preceding and 1 preceding), '') = '00000'"), # 确保前面没有额外的0:要么当前是第6行(前面正好5个记录),要么第6行前的event是1 F.or( F.col("row_num") == 6, F.lag("event", 6).over(user_time_window) == 1 ) ) ).otherwise(0) )
3. 提取每个用户的首个符合条件的1
如果一个用户有多个满足条件的1,我们只保留最早出现的那个:
# 给符合条件的行按用户排序,取第一个 target_df = df.filter(F.col("is_target") == 1).withColumn( "target_rank", F.row_number().over(Window.partitionBy("user_id").orderBy("month_col")) ).filter(F.col("target_rank") == 1).drop("row_num", "is_target", "target_rank")
为什么rank函数没用?
rank函数只能对分组内的行做排序标记,但你的需求需要识别连续的序列模式(恰好5个0+第一个1),这种状态依赖的逻辑,单纯的排序标记无法捕捉,必须用窗口函数跟踪前后行的状态或值。
边界情况说明
- 如果用户的前5条记录全是0,第6条是1:会被正确识别
- 如果用户有连续6个0后出现1:会被排除(因为前面不止5个0)
- 如果用户没有符合条件的行:不会出现在结果中,若需要保留所有用户,可以用左连接原表补充标记
内容的提问来源于stack exchange,提问作者user07
相关产品推荐
相关产品推荐

