You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 04:07:07