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

用Spark或SQL分析稀疏时间序列事件统计月度激活订阅用户数

稀疏状态变更事件的月末激活用户数统计方案

问题背景

你提到的这类场景属于时序状态快照统计,也被归为「缓慢变化维(Slowly Changing Dimension, SCD)类型2」的经典统计场景,专门处理随时间变化的离散状态的固定周期快照计算需求。


Spark 原生最优实现方案

核心依赖Spark内置窗口函数 last(expr, ignoreNulls=true),无需循环遍历时间窗口,也不需要提前补全冗余历史事件,全量一次性计算即可得到结果,分布式执行效率远高于你提到的两种方案。

具体实现步骤(Spark SQL 示例)

1. 预处理事件数据,保留用户每月最后一次状态

同一用户同一个月内的多次变更不会影响月末统计结果,只需保留当月最后一条事件即可:

WITH user_month_last_event AS (
    SELECT 
        user_id,
        event_type,
        date_trunc('month', time_stamp) AS event_month,
        -- 给状态打标记:激活为1,停用为0,方便后续聚合求和
        CASE WHEN event_type = 'A' THEN 1 ELSE 0 END AS status_flag
    FROM (
        SELECT 
            *,
            -- 按用户+月份分组,取当月时间最晚的一条事件
            ROW_NUMBER() OVER(PARTITION BY user_id, date_trunc('month', time_stamp) ORDER BY time_stamp DESC) AS rn
        FROM event_table
    ) t
    WHERE rn = 1
),

2. 生成统计周期内的连续月份序列

根据你的统计时间范围生成完整的月份列表,比如示例中统计2020年1月到5月:

all_months AS (
    SELECT explode(sequence(to_date('2020-01-01'), to_date('2020-05-01'), interval 1 month)) AS stat_month
)

3. 关联填充无事件月份的继承状态

将用户事件和所有统计月份关联后,用last函数向前填充无事件月份的状态:

user_month_all_status AS (
    SELECT 
        am.stat_month,
        u.user_id,
        -- 核心逻辑:按用户分区,按月份排序,取最近一次非空的状态值
        LAST(umle.status_flag, true) OVER(PARTITION BY u.user_id ORDER BY am.stat_month ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS final_status
    FROM all_months am
    -- 所有有过事件的用户和所有统计月份做笛卡尔关联
    CROSS JOIN (SELECT DISTINCT user_id FROM user_month_last_event) u
    -- 关联对应用户对应月份的事件状态
    LEFT JOIN user_month_last_event umle 
        ON u.user_id = umle.user_id 
        AND am.stat_month = umle.event_month
)

4. 按月份聚合得到最终激活用户数

SELECT 
    stat_month,
    SUM(final_status) AS active_user_count
FROM user_month_all_status
-- 过滤用户首次事件前的无效月份
WHERE final_status IS NOT NULL
GROUP BY stat_month
ORDER BY stat_month

方案优势

  1. 完全基于Spark内置函数实现,代码简洁易维护
  2. 无需循环遍历时间窗口,全量一次计算,分布式执行效率高
  3. 无需提前为用户补全所有月份的冗余事件,存储和计算成本都很低

如果你用DataFrame API开发,逻辑完全一致,对应org.apache.spark.sql.functions.last函数即可,第二个参数传true开启忽略空值的特性。

内容的提问来源于stack exchange,提问作者Robert

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 10:18:02