基于PySpark计算用户最长连续健身天数的聚合方法
问题描述
现有一张包含PersonID、Date、HasDoneWorkout字段的表,数据如下:
| PersonID | Date | HasDoneWorkout |
|---|---|---|
| A | 31-01-2001 | 1 |
| A | 01-02-2001 | 1 |
| A | 02-02-2001 | 1 |
| A | 03-02-2001 | 0 |
| A | 04-02-2001 | 1 |
| B | 02-02-2001 | 1 |
需求:创建PySpark聚合逻辑,统计每个用户连续健身的天数,若存在多段连续记录则取最长值。预期输出如下:
| PersonID | HasDoneWorkout |
|---|---|
| A | 3 |
| B | 1 |
此前尝试用Pandas方案但无法转换为PySpark实现,寻求解决方案。
PySpark解决方案
核心思路是通过分组排序+差值计算识别连续健身的分段,再统计各分段长度取最大值。具体实现步骤如下:
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("LongestContinuousWorkout").getOrCreate() # 模拟输入数据 data = [ ("A", "31-01-2001", 1), ("A", "01-02-2001", 1), ("A", "02-02-2001", 1), ("A", "03-02-2001", 0), ("A", "04-02-2001", 1), ("B", "02-02-2001", 1) ] df = spark.createDataFrame(data, ["PersonID", "Date", "HasDoneWorkout"]) # 1. 转换日期格式为PySpark日期类型 df = df.withColumn("Date", F.to_date("Date", "dd-MM-yyyy")) # 2. 定义窗口:按用户分组,按日期升序排序 user_window = Window.partitionBy("PersonID").orderBy("Date") # 3. 标记连续健身的分段 df = df.withColumn( "prev_date", F.lag("Date").over(user_window) ).withColumn( "date_diff", F.datediff("Date", "prev_date") ).withColumn( "segment_id", # 当健身状态为1时,若日期差不为1或为第一条记录,则生成新分段ID F.when( (F.col("HasDoneWorkout") == 1) & ((F.col("date_diff") != 1) | (F.col("date_diff").isNull())), F.monotonically_increasing_id() ).otherwise(F.lag("segment_id").over(user_window)) ) # 4. 统计每个分段的长度,取用户的最大值 result_df = df.filter(F.col("HasDoneWorkout") == 1) \ .groupBy("PersonID", "segment_id") \ .agg(F.count("*").alias("continuous_days")) \ .groupBy("PersonID") \ .agg(F.max("continuous_days").alias("HasDoneWorkout")) # 展示结果 result_df.show()
关键逻辑说明
- 日期转换:用
to_date指定格式dd-MM-yyyy将字符串转成日期类型,保证后续差值计算的准确性。 - 分段标记:通过
lag获取上一条记录的日期,计算日期差;当连续健身中断(日期差不为1或当天未健身)时,生成新的分段ID,以此区分不同的连续段。 - 聚合统计:先按用户和分段ID统计每个连续段的天数,再对每个用户取最大的分段长度,得到最终结果。
内容的提问来源于stack exchange,提问作者Krzysztof C
相关产品推荐
相关产品推荐

