Scala中基于薪资平均值映射新列(含字符串转整数及部门编码转换)
我来一步步帮你搞定这两个Scala Spark的DataFrame操作,话不多说直接上方案~
操作1:薪资列字符串转整数并生成基于平均值映射的新列
首先咱们得先把字符串类型的薪资转成整数,不然没法计算平均值。假设你的薪资列叫salary_str,先转成整数列salary;然后计算整个DataFrame的薪资平均值,再基于这个平均值给每行生成新列(这里给两种常见的映射示例,你可以根据实际需求调整)。
步骤分解:
- 字符串转整数:用
cast(IntegerType)或者Spark 3.0+的to_int函数,顺便处理可能的空值或非法值 - 计算薪资平均值:用聚合函数获取全局平均值,用广播变量优化性能避免重复计算
- 生成映射新列:自定义规则,比如标记是否高于平均值,或者计算薪资与平均值的比例
示例代码:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types.IntegerType // 假设原始DataFrame名为df val dfWithSalary = df // 字符串转整数,空值/非数字默认转0,你可以根据需求调整兜底值 .withColumn("salary", when(col("salary_str").cast(IntegerType).isNotNull, col("salary_str").cast(IntegerType)).otherwise(0)) // 获取全局薪资平均值,用广播变量减少数据传输 val avgSalary = dfWithSalary.agg(avg("salary")).first().getDouble(0) val broadcastAvg = spark.sparkContext.broadcast(avgSalary) // 示例1:生成标记列 - 1表示薪资高于平均值,0表示低于/等于 val dfWithAvgFlag = dfWithSalary .withColumn("is_above_avg", when(col("salary") > broadcastAvg.value, 1).otherwise(0)) // 示例2:生成比例列 - 薪资与平均值的比值(归一化映射) val dfWithAvgRatio = dfWithSalary .withColumn("salary_avg_ratio", col("salary") / broadcastAvg.value)
操作2:部门编码按平均薪资高低转换为数字
这个需求的核心是先按部门计算平均薪资,再给部门按平均薪资从高到低排序赋值数字(比如平均最高的部门对应1,次高对应2),最后把映射关联回原DataFrame替换部门编码。
步骤分解:
- 按部门分组计算平均薪资:用
groupBy+agg统计各部门薪资均值 - 给部门排序并分配排名:用窗口函数
row_number()或rank()(两者区别是处理相同均值的情况,按需选择) - 关联映射回原DataFrame:把排名数字替换原部门编码
示例代码:
import org.apache.spark.sql.expressions.Window // 先统计各部门平均薪资并分配排名 val deptSalaryRank = dfWithSalary .groupBy("dept_code") .agg(avg("salary").alias("dept_avg_salary")) .orderBy(desc("dept_avg_salary")) // 用row_number保证每个部门有唯一数字;若允许相同均值的部门同排名,换成rank() .withColumn("dept_rank", row_number().over(Window.orderBy(desc("dept_avg_salary")))) // 关联回原DataFrame,替换部门编码为排名数字 val dfWithDeptNum = dfWithSalary .join(deptSalaryRank, Seq("dept_code"), "left") .withColumn("dept_number", col("dept_rank")) // 不需要的中间列可以删掉,按需保留 .drop("dept_code", "dept_avg_salary", "dept_rank")
小提醒:
- 如果部门编码存在null值,记得用
coalesce给null部门设置默认编码,避免统计遗漏 - 若多个部门平均薪资相同,
rank()会给相同值分配相同排名,dense_rank()会跳过重复排名后的空位,根据业务需求选合适的窗口函数
内容的提问来源于stack exchange,提问作者Amey Totawar
相关产品推荐
相关产品推荐

