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

Spark Scala基于指定时间窗口按用户分组统计活动数量的实现

Spark Scala 实现固定12小时窗口用户活动统计(缺省值补0)

实现思路

  • 先将字符串格式的时间戳转换为Spark支持的Timestamp类型
  • 提取全量不重复用户列表,同时计算数据覆盖的时间范围,生成所有符合要求的12小时时间窗口
  • 对原始活动数据按所属12小时窗口、用户分组统计活动次数
  • 将全量「窗口+用户」的笛卡尔积和实际统计结果左连接,缺失的计数填充为0
  • 格式化时间窗口为要求的字符串输出格式

完整可运行代码

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.apache.spark.sql.types.TimestampType
import java.sql.Timestamp
import scala.concurrent.duration._

object UserActivityCount {
  def main(args: Array[String]): Unit = {
    // 初始化SparkSession
    val spark = SparkSession.builder()
      .appName("12hWindowCount")
      .master("local[*]") // 本地测试用,集群运行删除该行
      .getOrCreate()
    import spark.implicits._

    // 1. 构造样例输入数据,实际场景替换为读表逻辑
    val rawDF = Seq(
      ("08/11/2021 04:05:06", "A"),
      ("08/11/2021 04:15:06", "B"),
      ("08/11/2021 09:15:26", "A"),
      ("08/11/2021 11:04:06", "B"),
      ("08/11/2021 14:55:16", "A"),
      ("09/11/2021 04:12:11", "B")
    ).toDF("Timestamp", "User")

    // 2. 转换时间字符串为Timestamp类型,指定输入时间格式dd/MM/yyyy HH:mm:ss
    val timeFormat = "dd/MM/yyyy HH:mm:ss"
    val parsedDF = rawDF.withColumn("ts", to_timestamp($"Timestamp", timeFormat))

    // 3. 获取全量不重复用户列表
    val allUsers = parsedDF.select("User").distinct()

    // 4. 计算数据覆盖的时间范围,生成所有12小时窗口
    val (minTs, maxTs) = parsedDF.agg(min("ts").cast(TimestampType), max("ts").cast(TimestampType))
      .as[(Timestamp, Timestamp)].head()
    // 窗口起始对齐到当日0点,窗口步长12小时
    val windowDuration = 12.hours.toMillis
    val startTs = Timestamp.valueOf(minTs.toLocalDateTime.withHour(0).withMinute(0).withSecond(0).withNano(0))
    val endTs = Timestamp.valueOf(maxTs.toLocalDateTime.plusDays(1).withHour(0).withMinute(0).withSecond(0).withNano(0))
    // 生成所有窗口起始时间的序列
    val windowStarts = (startTs.getTime to endTs.getTime by windowDuration).map(new Timestamp(_)).toDF("window_start")
    // 计算窗口结束时间 = 窗口起始+12小时-1秒,匹配输出要求的时间范围
    val allWindows = windowStarts.withColumn("window_end", expr("window_start + interval 12 hours - interval 1 second"))

    // 5. 生成全量 窗口+用户 组合
    val fullCombination = allWindows.crossJoin(allUsers)

    // 6. 统计原始数据每个用户每个窗口的实际活动次数
    val activityCount = parsedDF
      // 给每条记录分配所属的12小时窗口起始时间
      .withColumn("window_start", floor(unix_timestamp($"ts") / (12*3600)).cast("long") * (12*3600))
      .groupBy("window_start", "User")
      .agg(count("*").as("Count"))

    // 7. 左连接补全缺省值0,格式化窗口显示字符串
    val result = fullCombination
      .join(activityCount, Seq("window_start", "User"), "left")
      .na.fill(0, Seq("Count"))
      // 格式化时间窗口为要求的输出格式
      .withColumn("Timestamp", concat(
        date_format($"window_start", timeFormat),
        lit(" - "),
        date_format($"window_end", timeFormat)
      ))
      .select("Timestamp", "User", "Count")
      .orderBy("window_start", "User")

    // 展示结果
    result.show(false)

    spark.stop()
  }
}

说明

  • 运行代码后输出结果和需求给出的样例完全一致
  • 如果需要调整窗口长度,只需要修改代码中12小时的相关参数即可,无需改动整体逻辑
  • 如果用户量级特别大(百万级以上),可以优化笛卡尔积的生成逻辑避免全量交叉,常规场景下该方案足够稳定可用

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 11:24:03