Spark/Spark SQL按streak分组计算连续活动日期的起止最值
问题说明
现有一张活动记录表,规则为:同一楼层下如果activity_date连续,则streak字段值递增,一旦日期不连续streak就重置为1。需要对每个连续的streak分组,获取每组活动日期的最小值与最大值,要求使用Spark+Scala或者Spark SQL实现。
输入表结构及样例数据
floor activity_date streak -------------------------------- floor1 2018-11-08 1 floor1 2019-01-24 1 floor1 2019-04-05 1 floor1 2019-04-08 1 floor1 2019-04-09 2 floor1 2019-04-14 1 floor1 2019-04-17 1 floor1 2019-04-20 1 floor2 2019-05-04 1 floor2 2019-05-05 2 floor2 2019-06-04 1 floor2 2019-07-28 1 floor2 2019-08-14 1 floor2 2019-08-22 1
预期输出结果
floor activity_date end_activity_date ---------------------------------------- floor1 2018-11-08 2018-11-08 floor1 2019-01-24 2019-01-24 floor1 2019-04-05 2019-04-05 floor1 2019-04-08 2019-04-09 floor1 2019-04-14 2019-04-14 floor1 2019-04-17 2019-04-17 floor1 2019-04-20 2019-04-20 floor2 2019-05-04 2019-05-05 floor2 2019-06-04 2019-06-04 floor2 2019-07-28 2019-07-28 floor2 2019-08-14 2019-08-14 floor2 2019-08-22 2019-08-22
实现方案
核心思路是通过标记连续streak的分组边界生成唯一分组ID,同一个分组内的日期属于同一个连续区间,最后分组聚合取首尾日期即可,两种实现方式逻辑完全一致:
1. Spark SQL实现
WITH -- 标记新分组起点:streak=1代表一个新连续区间的开始 group_flag AS ( SELECT floor, activity_date, CASE WHEN streak = 1 THEN 1 ELSE 0 END AS is_new_group FROM activity_table ), -- 累加标记生成分组ID,同一个连续区间的group_id相同 group_id AS ( SELECT floor, activity_date, SUM(is_new_group) OVER (PARTITION BY floor ORDER BY activity_date ASC) AS group_id FROM group_flag ) -- 按楼层和分组ID聚合,取每组最小、最大日期 SELECT floor, MIN(activity_date) AS activity_date, MAX(activity_date) AS end_activity_date FROM group_id GROUP BY floor, group_id ORDER BY floor, activity_date ASC;
2. Spark Scala DSL实现
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ // 定义窗口规则:按楼层分区,按活动日期升序排序 val windowSpec = Window.partitionBy("floor").orderBy("activity_date") val result = df // 生成分组起点标记 .withColumn("is_new_group", when(col("streak") === 1, 1).otherwise(0)) // 累加标记生成唯一分组ID .withColumn("group_id", sum("is_new_group").over(windowSpec)) // 分组聚合取首尾日期 .groupBy("floor", "group_id") .agg( min("activity_date").alias("activity_date"), max("activity_date").alias("end_activity_date") ) // 按要求排序后删除多余的分组ID字段 .orderBy("floor", "activity_date") .drop("group_id") // 输出结果 result.show(false)
内容的提问来源于stack exchange,提问作者SomeDataFellow
相关产品推荐
相关产品推荐

