Spark DataFrame groupBy+agg聚合后如何添加对应关联docId列
Spark分组取最大值同时获取关联字段实现方案
你当前的聚合逻辑仅计算了每个vocabId分组下count的最大值,没有保留最大值对应的docId信息,可以通过以下三种常用方案实现需求:
窗口函数方案(逻辑最直观)
按vocabId分区,对同分区内的数据按count倒序排序,取每个分区排名第1的行即可:import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ val windowSpec = Window.partitionBy("vocabId").orderBy(desc("count")) val result = docwords .withColumn("rank", row_number().over(windowSpec)) .filter($"rank" === 1) .select("docId", "vocabId", "count")如果同一个
vocabId下存在多个docId的count并列最大,row_number()会随机返回其中一条;需要返回所有并列结果可以替换为rank()函数。聚合后关联原表方案(易理解)
先按原有逻辑算出每个vocabId对应的最大count,再和原表做等值连接,匹配vocabId和count都相等的行:import org.apache.spark.sql.functions._ val maxCountDf = docwords.groupBy("vocabId").agg(max($"count") as "count") val result = docwords.join(maxCountDf, Seq("vocabId", "count")) .select("docId", "vocabId", "count")该方案会自动返回所有并列最大值的行,适合需要保留并列结果的场景。
struct聚合方案(性能最优)
Spark的max函数支持struct类型比较,比较时优先判断第一个字段值,再依次判断后续字段,因此可以把count和docId封装为struct取最大值,再拆分字段得到结果:import org.apache.spark.sql.functions._ val result = docwords .groupBy("vocabId") .agg(max(struct($"count", $"docId")) as "max_res") .select( $"max_res.docId" as "docId", $"vocabId", $"max_res.count" as "count" )
该方案仅需一次分组聚合shuffle,不需要额外join或窗口计算,大数据量下性能表现最好,并列最大值场景下仅返回其中一条。
三种方案输出结果均和你给出的预期结构一致,若需要固定结果顺序,可在最后追加orderBy($"vocabId")即可。
内容的提问来源于stack exchange,提问作者MattNewbie
相关产品推荐
相关产品推荐

