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

