用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
方案优势
- 完全基于Spark内置函数实现,代码简洁易维护
- 无需循环遍历时间窗口,全量一次计算,分布式执行效率高
- 无需提前为用户补全所有月份的冗余事件,存储和计算成本都很低
如果你用DataFrame API开发,逻辑完全一致,对应org.apache.spark.sql.functions.last函数即可,第二个参数传true开启忽略空值的特性。
内容的提问来源于stack exchange,提问作者Robert
相关产品推荐
相关产品推荐

