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

Spark Scala:分组时间窗口后提取含计数的最新时间戳记录

提取日级窗口内的最新最终计数行(Spark Scala实现)

嘿,这个需求我熟!咱们要做的就是对每个(user, app)的每日时间窗口,只保留该窗口内时间戳最晚的那一行——因为这行就是这个窗口的最终尝试次数,之前的都是中间累计的冗余数据对吧?下面直接上可落地的代码和逻辑说明:

步骤1:导入依赖包

首先得导入Spark SQL的函数和窗口函数相关包:

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

步骤2:处理时间列与生成日窗口标识

先把你的字符串格式时间转成Spark能识别的Timestamp类型,然后生成日级窗口的唯一标识(比如当日的起始时间,或者直接转成日期字符串):

// 假设你的原始DataFrame里,时间列叫event_time,计数列叫attempt_count
val dfWithTs = df
  .withColumn("event_ts", to_timestamp(col("event_time"), "yyyy-MM-dd HH:mm:ss")) // 转成Timestamp类型
  .withColumn("day_window", date_trunc("day", col("event_ts"))) // 生成日窗口标识(比如2017-12-22 00:00:00)

步骤3:用窗口函数筛选最新行

定义窗口分区规则:按user、app、day_window分组,然后在每个组内按时间戳降序排序,这样每个组里第一行就是时间最晚的记录:

// 定义窗口规格
val windowSpec = Window
  .partitionBy("user", "app", "day_window")
  .orderBy(col("event_ts").desc)

// 给每个组内的行打行号,最新的行号为1
val dfWithRowNum = dfWithTs.withColumn("row_num", row_number().over(windowSpec))

// 过滤出每个组的第一行,就是我们要的最终计数行
val finalDf = dfWithRowNum
  .filter(col("row_num") === 1)
  .drop("row_num", "day_window") // 不需要窗口标识的话可以删掉这两列

逻辑说明

  • 不管用户单日多次尝试(比如user1在2017-12-22有两次计数),还是跨多天尝试,这个逻辑都会精准提取每个(user, app, 日窗口)的最新时间戳记录,也就是该窗口的最终尝试次数。
  • 如果你的业务场景中存在同一窗口内时间戳完全相同的多条记录,可以把row_number()换成rank(),这样会保留所有时间戳相同的最新行;如果只需要任意一条,用row_number()就足够。

测试示例

比如你的初始数据是这样的:

val df = sc.parallelize(Seq(
  ("user1", "iphone", "2017-12-22 10:06:18", 1),
  ("user1", "iphone", "2017-12-22 14:30:22", 2),
  ("user1", "iphone", "2017-12-23 09:15:05", 1),
  ("user2", "android", "2017-12-22 11:20:33", 1)
)).toDF("user", "app", "event_time", "attempt_count")

运行上面的代码后,finalDf的结果会是:

userappevent_timeattempt_count
user1iphone2017-12-22 14:30:222
user1iphone2017-12-23 09:15:051
user2android2017-12-22 11:20:331

完全符合我们要的结果!

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 09:27:47