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

如何对DataFrame分组后按值计数排序并保留唯一项?

解决Spark DataFrame分组后按频率排序收集唯一值的问题

我来帮你搞定这个需求!你需要的是按col1分组后,获取col2的唯一值集合,并且按照每个值在分组内的出现次数降序排列对吧?之前用groupBy($"col1").count()确实只能统计每个col1分组的总条目数,没法满足细分到col2值的频率统计和排序需求,咱们可以分两步来实现:

步骤1:统计每个分组内col2各值的出现次数

首先我们需要先按col1和col2联合分组,统计每个组合的出现次数,这样就能得到每个col2值在对应col1分组里的频率:

Scala 版本

import org.apache.spark.sql.functions.{count, desc, collect_list, groupBy}

// 先统计(col1, col2)组合的出现次数
val countDF = df.groupBy($"col1", $"col2").agg(count("*").alias("cnt"))

PySpark 版本

from pyspark.sql import functions as F

# 统计每个(col1, col2)的出现次数
count_df = df.groupBy("col1", "col2").agg(F.count("*").alias("cnt"))

步骤2:分组排序并收集成列表

接下来我们按col1重新分组,然后在聚合时让col2的值按之前统计的cnt(次数)降序排列,最后收集成列表:

Scala 版本(Spark 2.4+支持)

val resultDF = countDF.groupBy($"col1")
  .agg(collect_list($"col2").orderBy(desc("cnt")).alias("col2_unique_sorted"))

PySpark 版本(Spark 2.4+支持)

result_df = count_df.groupBy("col1").agg(
    F.collect_list("col2").orderBy(F.desc("cnt")).alias("col2_unique_sorted")
)

兼容旧Spark版本的方案

如果你的Spark版本低于2.4,collect_list的orderBy语法可能不支持,那可以先对统计后的结果按col1和cnt降序排序,再分组收集:

Scala 版本

// 先按col1分组,再按次数降序排序
val sortedCountDF = countDF.orderBy($"col1", desc("cnt"))
// 分组收集列表
val resultDF = sortedCountDF.groupBy($"col1").agg(collect_list($"col2").alias("col2_unique_sorted"))

PySpark 版本

# 先按col1和cnt降序排序
sorted_count_df = count_df.orderBy("col1", F.desc("cnt"))
# 分组收集排序后的col2值
result_df = sorted_count_df.groupBy("col1").agg(F.collect_list("col2").alias("col2_unique_sorted"))

最终结果

运行完上面的代码后,你会得到符合预期的输出:

col1 | col2_unique_sorted
k1   | ["b", "a", "c"]
k2   | ["c", "b"]

内容的提问来源于stack exchange,提问作者NSR

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 07:50:41