PySpark实现按部门获取薪资最高与最低员工姓名
PySpark 分组获取各部门薪资最高/最低员工姓名
需求说明
现有员工部门、姓名、薪资数据,需按部门分组,输出每个部门薪资最高和最低的员工姓名。
输入数据
inputData = [(100,"ABC",2000),(100,"XYZ",1000),(100,"CDE",750),(200,"GYT",1500),(200,"JHU",1200),(200,"GHT",1300),(200,"YTR",8000)] inputSchema= "deptID int, empName String, empSal int" df = spark.createDataFrame(inputData,inputSchema) display(df)
解决方案
方法1:Spark 3.0+ 简洁实现(推荐)
Spark 3.0及以上版本提供了max_by和min_by函数,可直接在聚合操作中按指定字段排序,提取对应字段值:
from pyspark.sql import functions as F # 按部门分组,聚合获取薪资最高、最低的员工姓名 result_df = df.groupBy("deptID").agg( F.max_by("empName", "empSal").alias("maxSalEmp"), # 按薪资降序取对应员工姓名 F.min_by("empName", "empSal").alias("minSalEmp") # 按薪资升序取对应员工姓名 ) display(result_df)
方法2:兼容Spark低版本(窗口函数实现)
如果使用Spark 3.0以下版本,可通过窗口函数标记分组内的排序位次,再筛选聚合:
from pyspark.sql import Window from pyspark.sql import functions as F # 定义窗口:按部门分区,分别按薪资降序、升序排序 window_max = Window.partitionBy("deptID").orderBy(F.desc("empSal")) window_min = Window.partitionBy("deptID").orderBy(F.asc("empSal")) # 添加位次列,筛选位次为1的记录,再聚合得到结果 result_df = df.withColumn("rank_max", F.row_number().over(window_max)) \ .withColumn("rank_min", F.row_number().over(window_min)) \ .filter((F.col("rank_max") == 1) | (F.col("rank_min") == 1)) \ .groupBy("deptID") \ .agg( F.first(F.when(F.col("rank_max") == 1, F.col("empName"))).alias("maxSalEmp"), F.first(F.when(F.col("rank_min") == 1, F.col("empName"))).alias("minSalEmp") ) display(result_df)
输出结果
两种方法均可得到如下结果:
----------------------------- |deptID | maxSalEmp |minSalEmp| ----------------------------- | 100 | ABC | CDE | | 200 | YTR | JHU | -----------------------------
内容的提问来源于stack exchange,提问作者Vaibhav Gupta
相关产品推荐
相关产品推荐

