使用DataFrame的agg函数时如何保留其他列?附实例场景
如何在Spark DataFrame中筛选每个key对应Time最小的行并保留所有列
下面提供几种实用的解决方法:
方法1:窗口函数(推荐)
通过窗口函数对每个key分组并排序,标记行的排名后筛选目标行,能完整保留所有列,逻辑清晰且可控性强。
from pyspark.sql import Window import pyspark.sql.functions as F # 定义窗口:按key分组,按Time升序排序 window = Window.partitionBy("key").orderBy(F.col("Time").asc()) # 添加排名列,筛选排名为1的行后删除排名列 result_df = df.withColumn("row_rank", F.row_number().over(window)) \ .filter(F.col("row_rank") == 1) \ .drop("row_rank") result_df.show()
如果同一key下有多个行的Time值相同,row_number()会随机给它们分配不同排名;若想保留所有Time最小的行,可替换成F.rank()或F.dense_rank()。
方法2:分组聚合后关联原表
先提取每个key对应的最小Time,再通过关联操作匹配回原表中对应的整行。
import pyspark.sql.functions as F # 生成每个key对应的最小Time表 min_time_df = df.groupBy("key").agg(F.min("Time").alias("min_time")) # 关联原表,筛选出key和Time匹配的行 result_df = df.join( min_time_df, (df["key"] == min_time_df["key"]) & (df["Time"] == min_time_df["min_time"]), how="inner" ).drop(min_time_df["key"]) # 移除重复的key列 result_df.show()
方法3:排序后去重(简洁版)
利用dropDuplicates的特性,先按Time升序排序,再对key去重,保留每个key的第一行(即Time最小的行)。注意:如果同一key存在多个Time相同的行,此方法会随机保留其中一行。
result_df = df.orderBy("Time").dropDuplicates(["key"]) result_df.show()
内容的提问来源于stack exchange,提问作者syndromel
相关产品推荐
相关产品推荐

