PySpark计算等级首次出现至变更时的时间差总和
追踪等级变更的累计时间间隔计算方案
需求说明
需要计算用户等级(rank)从首次出现到发生变更时的累计时间差,例如用户6首次等级为rank2,后续变更为rank1时,累计时间差为该阶段所有时间差的总和(11.99+6.46=18.45)。
原始数据表
+-------+-------------------+-----------------+-------------------+-------------------+ |user_id| date_created| rank|date_diff_day_float| previous_rank| +-------+-------------------+-----------------+-------------------+-------------------+ | 1|2022-12-20 23:56:19| rank1| null| null| | 2|2022-12-26 19:28:50| rank1| null| null| | 2|2022-12-28 19:18:03| rank2| 1.99| rank1| | 2|2023-01-03 14:27:40| rank1| 5.8| rank2| | 3|2023-01-29 18:20:02| rank2| null| null| | 4|2023-01-10 18:27:59| rank1| null| null| | 5|2023-01-13 02:31:47| rank2| null| null| | 5|2023-01-23 23:17:40| rank2| 10.87| null| | 6|2022-12-22 18:49:35| rank2| null| null| | 6|2023-01-03 18:39:34| rank2| 11.99| null| | 6|2023-01-10 05:46:55| rank1| 6.46| rank2| | 6|2023-01-11 17:52:01| rank1| 1.5| null| | 6|2023-01-18 22:35:26| rank1| 7.2| null| | 6|2023-01-23 15:31:02| rank2| 4.71| rank1| | 7|2023-01-29 19:34:37| rank1| null| null| | 8|2023-03-06 18:38:19| rank1| null| null| | 9|2022-12-29 06:40:25| rank1| null| null| | 9|2023-01-25 17:27:09| rank1| 27.45| null| +-------+-------------------+-----------------+-------------------+-------------------+
期望输出
+-------+-------------------+-----------------+-------------------+-------------------+--------------------+ |user_id| date_created| rank|date_diff_day_float| previous_rank| rank_change| +-------+-------------------+-----------------+-------------------+-------------------+--------------------+ | 2|2022-12-28 19:18:03| rank2| 1.99| rank1| rank1=>rank2| | 2|2023-01-03 14:27:40| rank1| 5.8| rank2| rank2=>rank1| | 6|2023-01-10 05:46:55| rank1| 18.45| rank2| rank2=>rank1| | 6|2023-01-23 15:31:02| rank2| 13.41| rank1| rank1=>rank2| +-------+-------------------+-----------------+-------------------+-------------------+--------------------+
解决方案代码
基于PySpark实现,核心是通过窗口函数对用户的连续相同等级分组,再累计计算时间差:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 1. 按用户分组、时间排序,标记等级变更点并生成分组ID window_user = Window.partitionBy("user_id").orderBy("date_created") df = df.withColumn( "rank_change_flag", F.when(F.lag("rank").over(window_user) != F.col("rank"), 1).otherwise(0) ).withColumn( "rank_group", F.sum("rank_change_flag").over(window_user.rowsBetween(Window.unboundedPreceding, 0)) ) # 2. 在每个等级分组内,累计计算时间差总和(将null替换为0) window_group = Window.partitionBy("user_id", "rank_group").orderBy("date_created") df = df.withColumn( "cumulative_diff", F.sum(F.coalesce("date_diff_day_float", F.lit(0))).over(window_group) ) # 3. 筛选等级变更记录,替换时间差为累计值并生成变更描述 result_df = df.filter(F.col("previous_rank").isNotNull()) \ .withColumn("date_diff_day_float", F.col("cumulative_diff")) \ .withColumn( "rank_change", F.concat_ws("=>", F.col("previous_rank"), F.col("rank")) ) \ .select("user_id", "date_created", "rank", "date_diff_day_float", "previous_rank", "rank_change") # 查看结果 result_df.show()
代码说明
- 第一步通过
lag函数判断当前等级与上一条是否不同,生成变更标记,再累加标记得到每个连续等级的分组ID; - 第二步在每个分组内累计求和时间差,用
coalesce处理null值(首次记录的时间差为null,不影响累计); - 第三步筛选出有等级变更的记录,将累计时间差替换原字段,并生成变更描述文本。
内容的提问来源于stack exchange,提问作者Pedro Daumas
相关产品推荐
相关产品推荐

