如何使用PySpark统计从指定日期倒推的用户连续登录天数
基于PySpark计算指定日期倒推的用户连续登录天数
需求说明
从指定日期2022-01-04倒推,统计各用户在应用中的连续登录天数,因数据规模较大,采用PySpark实现。
输入数据
Name Date John 2022-01-01 John 2022-01-01 Mary 2022-01-01 Steve 2022-01-03 Mary 2022-01-03 John 2022-01-02 John 2022-01-03 Mary 2022-01-04 John 2022-01-04
期望输出
Name consecutive_days John 4 Mary 2
实现步骤与代码
核心思路
- 去重处理:同一用户同一天的重复登录记录仅保留一条,避免重复计算
- 日期筛选:仅保留
2022-01-04及之前的登录数据 - 分组排序:按用户分组,将登录日期按降序排列
- 连续日期分组:通过计算日期与指定日期的间隔天数,结合行号生成分组标识,连续登录的日期会得到相同的分组键
- 统计连续天数:筛选出从指定日期开始的连续分组,统计该分组内的记录数即为连续登录天数
PySpark代码实现
from pyspark.sql import SparkSession from pyspark.sql.functions import col, date_diff, row_number from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("ConsecutiveLoginDays").getOrCreate() # 读取输入数据(实际场景可替换为读取数据库或分布式文件) data = [ ("John", "2022-01-01"), ("John", "2022-01-01"), ("Mary", "2022-01-01"), ("Steve", "2022-01-03"), ("Mary", "2022-01-03"), ("John", "2022-01-02"), ("John", "2022-01-03"), ("Mary", "2022-01-04"), ("John", "2022-01-04") ] df = spark.createDataFrame(data, ["Name", "Date"]) # 将字符串日期转换为日期类型 df = df.withColumn("Date", col("Date").cast("date")) # 1. 去重:同一用户同一天只保留一条记录 df_distinct = df.dropDuplicates(["Name", "Date"]) # 2. 筛选指定日期及之前的记录 target_date = "2022-01-04" df_filtered = df_distinct.filter(col("Date") <= target_date) # 3. 按用户分组,日期降序排列,生成行号(目标日期对应行号1) window_spec = Window.partitionBy("Name").orderBy(col("Date").desc()) df_ranked = df_filtered.withColumn("rn", row_number().over(window_spec)) # 4. 生成连续登录分组键:连续日期的group_key值相同 df_grouped = df_ranked.withColumn( "group_key", date_diff(col("Date"), target_date) + col("rn") ) # 5. 筛选包含目标日期的分组,统计组内记录数即为连续天数 result = df_grouped.filter(col("Date") == target_date)\ .join(df_grouped, on=["Name", "group_key"], how="inner")\ .groupBy("Name")\ .agg(col("rn").max().alias("consecutive_days"))\ .orderBy("Name") # 输出结果 result.show()
代码关键逻辑解释
- 去重:
dropDuplicates(["Name", "Date"])避免同一用户同一天的多次登录被重复计数 - 窗口排序:按用户分组后降序排列日期,行号
rn从1开始(对应目标日期) - 分组键生成:
date_diff(Date, target_date)计算日期与目标日期的间隔(目标日期为0,前一天为-1),加上行号后,连续日期的group_key值完全相同,非连续日期会产生不同的键 - 统计连续天数:通过关联包含目标日期的分组,统计组内最大行号,即为从目标日期倒推的连续登录天数
结果验证
运行代码后输出与期望结果一致:
+----+----------------+ |Name|consecutive_days| +----+----------------+ |John| 4| |Mary| 2| +----+----------------+
内容的提问来源于stack exchange,提问作者karek77
相关产品推荐
相关产品推荐

