Spark DataFrame基于优先级层级的无重叠Pivot实现方案问询
适配任意优先级层级的实现方案
核心思路是先对同一个分组内的每一个Column_D取值,仅保留其最高优先级的对应记录,再执行pivot操作,从根源上避免重复值出现在低优先级列。
具体实现步骤
- 定义优先级映射规则
按照业务给定的优先级分数定义映射关系,分数越高代表优先级越高,可灵活调整优先级数量:
// 可根据业务需求任意修改优先级数量和分数 val priorityScoreMap = Map( "very high" -> 1.0, "high" -> 0.8, "medium" -> 0.6, "low" -> 0.4, "very low" -> 0.2 ) // 示例也可替换为你给出的短优先级配置:Map("high"->1, "medium"->0.5, "low"->0.25)
- 给原始数据增加优先级分数字段
通过UDF将Priority文本映射为可排序的数值:
import org.apache.spark.sql.functions.udf // 可按需增加容错逻辑,未定义的优先级默认按最低分处理 val getPriorityScore = udf((priority: String) => priorityScoreMap.getOrElse(priority, 0.0)) val dfWithScore = df.withColumn("priority_score", getPriorityScore($"Priority"))
- 对Column_D去重,仅保留最高优先级记录
用窗口函数筛选每个Column_A、Column_B、Column_C、Column_D组合下优先级最高的唯一记录:
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions.row_number val winSpec = Window.partitionBy("Column_A", "Column_B", "Column_C", "Column_D") .orderBy($"priority_score".desc) // 筛选后每个D值仅会出现在其最高优先级的记录中 val dfDedup = dfWithScore.withColumn("rn", row_number().over(winSpec)) .filter($"rn" === 1) .drop("rn", "priority_score")
- 执行pivot得到预期结果
自动按照优先级从高到低排列pivot后的列:
// 提取按优先级从高到低排序的列名,保证pivot后列顺序符合预期 val priorityOrder = priorityScoreMap.toSeq.sortBy(-_._2).map(_._1) val outcome = dfDedup.groupBy("Column_A", "Column_B", "Column_C") .pivot("Priority", priorityOrder) .agg(collect_set("Column_D") as "Set_D")
方案优势
- 完全适配任意数量的优先级层级,仅需修改
priorityScoreMap即可,无需调整核心逻辑 - 性能稳定,无论优先级层级多少,都仅需一次去重+一次pivot操作,无额外的集合运算开销
- 自动适配优先级列顺序,无需手动指定pivot列的排序规则
内容的提问来源于stack exchange,提问作者Shivam Chaudhary
相关产品推荐
相关产品推荐

