PySpark计算客户连续活跃月数:重置累积和需保留前值
PySpark 计算客户连续活跃月数优化方案
核心思路
针对「非活跃首行保留累计值、后续非活跃行归0」的需求,通过窗口函数分层计算实现,无需额外分组后二次处理,逻辑简洁高效。
实现代码
from pyspark.sql import Window from pyspark.sql.functions import col, lag, sum, when, coalesce # 1. 定义客户级窗口:按客户分区,按月初日期排序 customer_window = Window.partitionBy("customerId").orderBy("month_first_date") # 2. 生成活跃分组标识:客户从非活跃转为活跃时,分组ID递增 df = df.withColumn( "group_id", sum( when( col("active_in_month") & (coalesce(lag("active_in_month").over(customer_window), False) == False), 1 ).otherwise(0) ).over(customer_window) ) # 3. 定义分组级窗口:按客户+活跃分组分区,按月初日期排序 group_window = Window.partitionBy("customerId", "group_id").orderBy("month_first_date") # 4. 计算分组内的活跃月累计数 df = df.withColumn( "active_cum_sum", sum(col("active_in_month").cast("int")).over(group_window) ) # 5. 生成目标列:匹配需求规则 df = df.withColumn( "ConsecMonthsActiveTarget", when( col("active_in_month"), col("active_cum_sum") ).when( ~col("active_in_month") & col("active_previous_month"), lag("active_cum_sum").over(customer_window) ).otherwise(0) )
方案优势
- 无冗余步骤:直接通过窗口函数完成分组、累计、规则判断,避免先计算
CumSumWithGroups再二次修正的繁琐流程 - 性能更优:所有操作基于窗口函数完成,减少不必要的shuffle开销
- 逻辑清晰:每一步对应需求规则,便于维护和调试
内容的提问来源于stack exchange,提问作者Slite
相关产品推荐
相关产品推荐

