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的结果会是:
| user | app | event_time | attempt_count |
|---|---|---|---|
| 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 |
完全符合我们要的结果!
内容的提问来源于stack exchange,提问作者annonymous_guy
相关产品推荐
相关产品推荐

