如何对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
相关产品推荐
相关产品推荐

