Spark/Scala如何通过另一表映射对DataFrame中的id列表排序
实现方案
Spark Scala环境下有两种常用实现方式,可根据DF2的体量选择:
方案1:广播小表+UDF(适合DF2数据量小的场景)
核心思路是把权重映射关系广播到所有Executor,用自定义UDF逐行对id列表按权重排序,性能损耗极低。
代码实现
import org.apache.spark.sql.functions.col import org.apache.spark.sql.functions.udf // 假设DF1结构:pk:Int, id_list:Array[String] // 假设DF2结构:id:String, weight:Int // 1. 收集DF2的id-权重映射并广播 val weightMap = df2.select("id", "weight") .collect() .map(row => row.getAs[String]("id").trim -> row.getAs[Int]("weight")) .toMap val broadcastWeights = spark.sparkContext.broadcast(weightMap) // 2. 定义排序UDF val sortByWeightUdf = udf { (idList: Seq[String]) => val weights = broadcastWeights.value // 按权重升序排列,不存在权重的id默认排在末尾,可自定义默认值调整顺序 idList.sortBy(id => weights.getOrElse(id, Int.MaxValue)) } // 3. 生成结果 val resultDf = df1.withColumn("sorted_id_list", sortByWeightUdf(col("id_list")))
方案2:内置函数实现(适合DF2数据量大无法广播的场景)
核心思路是先炸开id列表、关联权重表、组内排序后再聚合回列表,全程用Spark内置函数,不需要广播变量。
代码实现
import org.apache.spark.sql.functions._ val resultDf = df1 // 把每行的id列表拆成单条id行 .select(col("pk"), explode(col("id_list")).alias("id")) // 左关联权重表,避免id不存在时丢失数据 .join(df2, Seq("id"), "left") // 按主键分组,组内按权重升序排序后重新收集为列表 .groupBy("pk") .agg( collect_list(col("id")) .within_group(orderBy(col("weight").asc_nulls_last)) .alias("sorted_id_list") )
注意事项
你示例中DF2的kangaro存在拼写错误(缺末尾的o),实际运行前需要先对齐两个表的id取值,否则关联不到权重会导致排序错位。如果需要把无权重的id排在列表开头,把asc_nulls_last改为asc_nulls_first即可。
内容的提问来源于stack exchange,提问作者219CID
相关产品推荐
相关产品推荐

