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

如何使用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 

实现步骤与代码

核心思路

  1. 去重处理:同一用户同一天的重复登录记录仅保留一条,避免重复计算
  2. 日期筛选:仅保留2022-01-04及之前的登录数据
  3. 分组排序:按用户分组,将登录日期按降序排列
  4. 连续日期分组:通过计算日期与指定日期的间隔天数,结合行号生成分组标识,连续登录的日期会得到相同的分组键
  5. 统计连续天数:筛选出从指定日期开始的连续分组,统计该分组内的记录数即为连续登录天数

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 08:36:10