如何高效(近似)获取Spark DataFrame分组内的最高出现频次
Spark分组内获取最高频次记录的实现方案
需求说明
基于给定的Spark DataFrame,按ColA分组后,获取每组内ColB出现频次最高的记录。数据集规模为9.5亿条记录、20400个分区,允许近似结果,但禁止使用first()/last()方法。
实现方案
方案一:精确统计(适合对结果准确性要求高的场景)
通过分组统计+窗口函数实现,步骤如下:
- 先统计
ColA与ColB组合的出现频次; - 按
ColA分组,对频次降序排序后,用行号函数筛选每组的第一条记录。
from pyspark.sql import Window from pyspark.sql.functions import desc, row_number # 1. 统计(ColA, ColB)组合的频次 count_df = df.groupBy("ColA", "ColB").count() # 2. 定义窗口规则:按ColA分组,按count降序排列 window_spec = Window.partitionBy("ColA").orderBy(desc("count")) # 3. 筛选每组频次最高的记录 result_df = count_df.withColumn("row_num", row_number().over(window_spec)) \ .filter("row_num == 1") \ .drop("row_num") result_df.show()
性能优化提示:针对超大数据集,可保持count_df的分区数与原数据集一致(如count_df.repartition(20400)),减少shuffle开销。
方案二:近似统计(适合超大数据集,性能优先)
使用Spark内置的freqItems近似统计函数,基于概率算法获取分组内的频繁项,无需全量shuffle,性能更优。
from pyspark.sql.functions import freqItems, explode, col # 1. 获取每个ColA分组中ColB的近似频繁项(参数0.1为支持度阈值,可调整) freq_df = df.groupBy("ColA").agg(freqItems(["ColB"], 0.1).alias("freq_ColB")) # 2. 展开频繁项列并关联频次统计结果 count_df = df.groupBy("ColA", "ColB").count() result_df = freq_df.select("ColA", explode(col("freq_ColB")).alias("ColB")) \ .join(count_df, ["ColA", "ColB"], "inner") \ .withColumn("row_num", row_number().over(Window.partitionBy("ColA").orderBy(desc("count")))) \ .filter("row_num == 1") \ .drop("row_num") result_df.show()
说明:freqItems的第二个参数为支持度阈值,值越小结果越精确但性能稍有下降;该方法返回的是分组内出现频率较高的项集合,结合后续的频次排序可稳定获取最高频记录,避免了first()/last()的不确定性。
内容的提问来源于stack exchange,提问作者Lelouch Lamperouge
相关产品推荐
相关产品推荐

