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

基于条件重置计数的DataFrame连续活跃月窗口函数实现需求

解决连续活跃月份计数(窗口函数实现)

核心思路

要实现间隔超1个月重置计数的需求,关键是先把同一用户的连续活跃时间段划分成独立分组,再在每个分组内做连续计数。单纯用无边界窗口无法处理重置,必须先通过月份差标记新分组,再累加得到分组ID,最后在分组内计数。

具体实现(PySpark为例)

假设你的DataFrame名为user_active_df,包含user_id(用户ID)、reference_month(用户活跃月份,格式为日期类型,比如yyyy-MM)。

  1. 计算当前行与上一行的月份间隔
    先按用户分区、活跃月份排序,用lag函数取上一行的活跃月份,再计算月份差:

    from pyspark.sql import functions as F
    from pyspark.sql.window import Window
    
    # 定义按用户分区、按活跃月份排序的窗口
    user_order_window = Window.partitionBy("user_id").orderBy("reference_month")
    
    user_active_df = user_active_df.withColumn(
        "prev_active_month",
        F.lag("reference_month").over(user_order_window)
    ).withColumn(
        "month_interval",
        F.months_between("reference_month", "prev_active_month")
    )
    
  2. 标记新分组并生成分组ID
    当是用户的第一条记录,或者当前与上一次活跃间隔超过1个月时,标记为新分组;然后累加标记值得到每个连续活跃段的分组ID:

    user_active_df = user_active_df.withColumn(
        "is_new_continuous_group",
        F.when(
            F.col("prev_active_month").isNull() | (F.col("month_interval") > 1),
            1
        ).otherwise(0)
    ).withColumn(
        "continuous_group_id",
        F.sum("is_new_continuous_group").over(user_order_window)
    )
    
  3. 在分组内计算连续活跃月份数
    按用户+分组ID分区、活跃月份排序,用row_number()生成连续计数:

    group_order_window = Window.partitionBy("user_id", "continuous_group_id").orderBy("reference_month")
    
    user_active_df = user_active_df.withColumn(
        "active_months",
        F.row_number().over(group_order_window)
    )
    

关键说明

  • 这里用months_between计算月份差,确保跨年度的间隔也能正确判断(比如2023-12到2024-01的间隔是1,不会误判为重置)。
  • 如果你用的是Spark SQL,逻辑完全一致,只是语法换成SQL写法:
    WITH step1 AS (
        SELECT 
            user_id,
            reference_month,
            LAG(reference_month) OVER(PARTITION BY user_id ORDER BY reference_month) AS prev_active_month,
            MONTHS_BETWEEN(reference_month, prev_active_month) AS month_interval
        FROM user_active_df
    ),
    step2 AS (
        SELECT 
            *,
            CASE 
                WHEN prev_active_month IS NULL OR month_interval > 1 THEN 1 
                ELSE 0 
            END AS is_new_continuous_group,
            SUM(is_new_continuous_group) OVER(PARTITION BY user_id ORDER BY reference_month) AS continuous_group_id
        FROM step1
    )
    SELECT 
        *,
        ROW_NUMBER() OVER(PARTITION BY user_id, continuous_group_id ORDER BY reference_month) AS active_months
    FROM step2;
    

内容的提问来源于stack exchange,提问作者Thais Guerra Braga

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 00:56:16