基于条件重置计数的DataFrame连续活跃月窗口函数实现需求
解决连续活跃月份计数(窗口函数实现)
核心思路
要实现间隔超1个月重置计数的需求,关键是先把同一用户的连续活跃时间段划分成独立分组,再在每个分组内做连续计数。单纯用无边界窗口无法处理重置,必须先通过月份差标记新分组,再累加得到分组ID,最后在分组内计数。
具体实现(PySpark为例)
假设你的DataFrame名为user_active_df,包含user_id(用户ID)、reference_month(用户活跃月份,格式为日期类型,比如yyyy-MM)。
计算当前行与上一行的月份间隔
先按用户分区、活跃月份排序,用lag函数取上一行的活跃月份,再计算月份差:from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义按用户分区、按活跃月份排序的窗口 user_order_window = Window.partitionBy("user_id").orderBy("reference_month") user_active_df = user_active_df.withColumn( "prev_active_month", F.lag("reference_month").over(user_order_window) ).withColumn( "month_interval", F.months_between("reference_month", "prev_active_month") )标记新分组并生成分组ID
当是用户的第一条记录,或者当前与上一次活跃间隔超过1个月时,标记为新分组;然后累加标记值得到每个连续活跃段的分组ID:user_active_df = user_active_df.withColumn( "is_new_continuous_group", F.when( F.col("prev_active_month").isNull() | (F.col("month_interval") > 1), 1 ).otherwise(0) ).withColumn( "continuous_group_id", F.sum("is_new_continuous_group").over(user_order_window) )在分组内计算连续活跃月份数
按用户+分组ID分区、活跃月份排序,用row_number()生成连续计数:group_order_window = Window.partitionBy("user_id", "continuous_group_id").orderBy("reference_month") user_active_df = user_active_df.withColumn( "active_months", F.row_number().over(group_order_window) )
关键说明
- 这里用
months_between计算月份差,确保跨年度的间隔也能正确判断(比如2023-12到2024-01的间隔是1,不会误判为重置)。 - 如果你用的是Spark SQL,逻辑完全一致,只是语法换成SQL写法:
WITH step1 AS ( SELECT user_id, reference_month, LAG(reference_month) OVER(PARTITION BY user_id ORDER BY reference_month) AS prev_active_month, MONTHS_BETWEEN(reference_month, prev_active_month) AS month_interval FROM user_active_df ), step2 AS ( SELECT *, CASE WHEN prev_active_month IS NULL OR month_interval > 1 THEN 1 ELSE 0 END AS is_new_continuous_group, SUM(is_new_continuous_group) OVER(PARTITION BY user_id ORDER BY reference_month) AS continuous_group_id FROM step1 ) SELECT *, ROW_NUMBER() OVER(PARTITION BY user_id, continuous_group_id ORDER BY reference_month) AS active_months FROM step2;
内容的提问来源于stack exchange,提问作者Thais Guerra Braga
相关产品推荐
相关产品推荐

