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
相关产品推荐
相关产品推荐

