Spark Window函数使用问题求助:无法正确提取最高薪资
问题原因分析
你的问题核心出在Window函数的分区逻辑上,咱们一步步拆解:
- 最初用
partitionBy("id"):你的数据集里每个id都是唯一的(从0到29每条记录对应一个独立id),这意味着每个分区里只有1条数据。rank()函数是在每个分区内计算排名,单条数据的排名自然都是1,根本没法跨记录比较薪资高低。 - 后来改成按
Salary分区:你的数据里几乎每个薪资值都是唯一的,所以每个分区还是只有1条数据,结果和之前完全一样,所有rank值都是1。
解决方案
要找出全局最高薪资,咱们不需要对数据分区,直接在整个数据集范围内排序取排名第一的记录就行,有两种实用方法:
方法1:使用窗口函数(适合需要保留排名信息的场景)
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions.{rank, col} // 定义全局窗口,按Salary降序排序(不需要partitionBy,代表整个数据集为一个窗口) val globalWindow = Window.orderBy(col("Salary").desc) // 为每条记录添加薪资排名列 val rankedDF = df1.withColumn("salary_rank", rank().over(globalWindow)) // 筛选出排名第一的记录 val highestSalaryDF = rankedDF.filter(col("salary_rank") === 1) highestSalaryDF.show()
执行后你会看到只有薪资122391.0的那条记录,salary_rank值为1。
方法2:直接排序取第一条(更高效,无需额外列)
如果只需要最高薪资的记录,不需要保留排名信息,直接用orderBy加limit会更简单高效:
val highestSalaryDF = df1.orderBy(col("Salary").desc).limit(1) highestSalaryDF.show()
这个方法会直接返回薪资最高的那条记录,也就是你要的122391.0对应的行。
内容的提问来源于stack exchange,提问作者Amal Swain
相关产品推荐
相关产品推荐

