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

Scala Spark中DataFrame按小时聚合及evtTime小时提取需求问询

Scala Spark 实现数据按小时聚合方案

没问题,我来帮你搞定这个按小时聚合数据的需求!咱们一步步来操作:

1. 先处理时间字段,提取小时维度

首先需要把evtTime字段转换为日期+小时的格式(比如2018-01-01 11),这样就能精准按小时维度分组。这里分两种常见情况处理:

情况1:evtTime已经是Timestamp类型

直接用date_format函数格式化即可:

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

val dfWithHour = df.withColumn("hourly_time", date_format(col("evtTime"), "yyyy-MM-dd HH"))

情况2:evtTime是字符串类型

先把字符串转成Timestamp类型,再进行格式化:

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

val dfWithHour = df
  .withColumn("evtTime_ts", to_timestamp(col("evtTime"), "yyyy-MM-dd HH:mm:ss.SSS"))
  .withColumn("hourly_time", date_format(col("evtTime_ts"), "yyyy-MM-dd HH"))
  .drop("evtTime_ts") // 不需要临时字段的话可以删掉

2. 按用户+小时维度聚合数据

接下来就可以按reqUser和刚生成的hourly_time分组,对event_count进行聚合(这里用求和示例,你可以根据需求换成计数、平均值等):

val hourlyAggDF = dfWithHour
  .groupBy("reqUser", "hourly_time")
  .agg(sum("event_count").alias("total_event_count"))
  .orderBy("reqUser", "hourly_time") // 可选,按用户和时间排序让结果更清晰

3. 查看聚合结果

执行hourlyAggDF.show()就能看到最终的聚合结果,你的示例数据会输出类似这样的内容:

+-------+-------------+------------------+
|reqUser| hourly_time |total_event_count|
+-------+-------------+------------------+
|X166814|2018-01-01 11|                 1|
|X166815|2018-01-01 02|                 2|
|X166816|2018-01-01 11|                 5|
|X166817|2018-02-01 10|                 1|
|X166818|2018-01-01 09|                 3|
|X166819|2018-01-01 10|                 8|
+-------+-------------+------------------+

如果更倾向于保留Timestamp类型的小时维度(比如2018-01-01 11:00:00),也可以用trunc函数直接截断时间到小时:

val dfWithHourTrunc = df.withColumn("hourly_time", trunc(col("evtTime"), "hour"))

后续的聚合逻辑和上面完全一致。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 06:48:12