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

Scala Spark DataFrame:如何保留同Key最小Time行并删除其余行

Scala Spark 保留每个Key对应最小Time的行

核心注意点

你的Time字段是StringType,直接用字符串比较会出现逻辑错误(比如字符串"10"会被认为比"2"小),所以必须先将其转换为数值类型(如IntegerType)再进行最小值判断。

方法一:窗口函数(推荐)

窗口函数可以直接对每个key分组后排序,标记出每组中Time最小的行,再过滤保留目标行。

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

// 1. 定义窗口规则:按key分组,按转换为Int的Time升序排序
val windowSpec = Window.partitionBy("key").orderBy(col("Time").cast("Int").asc)

// 2. 给每行添加行号,每组内Time最小的行号为1
val rankedDF = df.withColumn("row_num", row_number().over(windowSpec))

// 3. 过滤出行号为1的行,再删除辅助列
val resultDF = rankedDF.filter(col("row_num") === 1).drop("row_num")

resultDF.show()

方法二:GroupBy + Join

先通过分组聚合找到每个key对应的最小Time,再和原表关联筛选目标行。

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

// 1. 计算每个key的最小Time(转换为Int计算,再转回String保持原类型)
val minTimeDF = df.groupBy("key")
                  .agg(min(col("Time").cast("Int")).cast("String").alias("min_Time"))

// 2. 和原表关联,保留key和Time匹配的行
val resultDF = df.join(minTimeDF, Seq("key"), "inner")
                 .filter(col("Time") === col("min_Time"))
                 .drop("min_Time")

resultDF.show()

为什么drop/filter无法实现?

  • dropDuplicates只能删除完全相同的行,无法针对分组保留特定条件的行;
  • filter是逐行判断,无法跨分组比较并筛选出每组的最小值行,必须结合分组或窗口逻辑才能实现需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 07:25:17