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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 02:54:05