You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.03 03:57:37