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

