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

Scala Spark按时间桶+KEY分组统计跨桶事件duration总和方案咨询

核心实现思路
  • 先将事件的开始时间转换为Epoch秒格式,同时计算事件的结束时间 = 开始时间 + 耗时duration
  • 按5分钟(300秒)的桶粒度,生成每个事件覆盖的所有时间桶列表,通过explode拆分为多行
  • 逐行计算事件在对应桶内的有效时长:有效时长 = min(事件结束时间, 桶结束时间) - max(事件开始时间, 桶开始时间)
  • 最后按KEY和桶时间分组,求和有效时长即可得到符合要求的统计结果

推荐直接使用Epoch时间戳格式计算,不需要来回做时间格式转换,执行效率更高。


代码实现示例

Scala DataFrame API 版本

import org.apache.spark.sql.functions._
import org.apache.spark.sql.types.LongType

// 步骤1:预处理生成事件的开始、结束时间戳
val processedDF = df
  // 如果Time字段已经是Epoch秒格式,直接重命名为start_ts即可,无需调用unix_timestamp转换
  .withColumn("start_ts", unix_timestamp($"Time", "yyyy-MM-dd HH:mm:ss").cast(LongType))
  .withColumn("end_ts", $"start_ts" + $"duration")

// 步骤2:拆分跨桶事件为多行,计算每个桶的有效时长
val splitDF = processedDF
  // 生成事件覆盖的所有5分钟桶的起始时间戳
  .withColumn("bucket_start", explode(sequence(
    floor($"start_ts" / 300) * 300,
    floor($"end_ts" / 300) * 300,
    lit(300)
  )))
  .withColumn("bucket_end", $"bucket_start" + 300)
  .withColumn("valid_duration", least($"end_ts", $"bucket_end") - greatest($"start_ts", $"bucket_start"))

// 步骤3:分组聚合得到结果
val resultDF = splitDF
  .groupBy($"KEY", $"bucket_start")
  .agg(sum($"valid_duration").alias("total_duration"))
  // 可选:把桶时间转回可读字符串格式
  .withColumn("bucket_time", from_unixtime($"bucket_start", "yyyy-MM-dd HH:mm:ss"))
  .select("KEY", "bucket_time", "total_duration")

Spark SQL 版本

WITH processed_event AS (
  SELECT
    KEY,
    -- 如果Time字段已经是Epoch秒格式,直接写 Time AS start_ts 即可
    unix_timestamp(Time, 'yyyy-MM-dd HH:mm:ss') AS start_ts,
    unix_timestamp(Time, 'yyyy-MM-dd HH:mm:ss') + duration AS end_ts
  FROM your_table_name
),
event_bucket AS (
  SELECT
    KEY,
    bucket_start,
    LEAST(end_ts, bucket_start + 300) - GREATEST(start_ts, bucket_start) AS valid_duration
  FROM processed_event
  LATERAL VIEW EXPLODE(SEQUENCE(FLOOR(start_ts/300)*300, FLOOR(end_ts/300)*300, 300)) t AS bucket_start
)
SELECT
  KEY,
  FROM_UNIXTIME(bucket_start, 'yyyy-MM-dd HH:mm:ss') AS bucket_time,
  SUM(valid_duration) AS total_duration
FROM event_bucket
GROUP BY KEY, bucket_start

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 14:15:05