公司动态

Spark 核心之二次排序、分组取 TopN 优化分析

📅 2026/8/14 14:18:14
Spark 核心之二次排序、分组取 TopN 优化分析
摘要为什么你的分组 TopN 任务跑了 30 分钟还 OOM答案通常是四个字——groupByKey。本文将二次排序和分组 TopN 作为两条优化主线从反模式剖析、三种方案对比groupByKey / reduceByKeymapPartitions / SQL 窗口函数、groupByKey vs reduceByKey 本质区别、完整代码实战、Spark SQL 窗口函数最佳实践五个维度配合 1 张原创深色架构图助你彻底攻克 TopN 优化的每一个细节。关键词Spark 二次排序, TopN, groupByKey, reduceByKey, 窗口函数, row_number, mapPartitions, 分组优化一、开篇groupByKey 是 TopN 的超级陷阱// ❌ 这是 95% 的新手会写的代码——也是 OOM 的根源rdd.groupByKey()// 全量数据按 Key 分组.mapValues(_.toList.sortBy(-_._2).take(3))// ⚠ 所有 Value 加载到内存// 问题: 如果某个 Key 有 1000 万条数据 → Executor OOMgroupByKey 的问题它不做 Map 端预聚合将全部 KV 对原封不动地 Shuffle 到 Reducer然后全部加载到内存——数据倾斜时直接炸掉。二、优化方案全景图三、方案 1二次排序Secondary Sort3.1 自定义排序键// 按 classId 升序score 降序caseclassSecondarySortKey(classId:String,score:Int)extendsOrdered[SecondarySortKey]{overridedefcompare(that:SecondarySortKey):Int{valcmp1this.classId.compareTo(that.classId)if(cmp1!0)cmp1else-this.score.compareTo(that.score)// score 降序}}// 使用rdd.map{case(classId,name,score)(SecondarySortKey(classId,score),(name,score))}.sortByKey()// 按自定义键全局排序3.2 逐组处理避免 groupByKey// 利用 sortByKey 后的有序性mapPartitions 内逐组处理rdd.mapPartitions{itervarcurrentClassvarbufferListBuffer[(String,Int)]()iter.flatMap{case(key,(name,score))if(key.classId!currentClass){valresultbuffer.toList currentClasskey.classId bufferListBuffer((name,score))if(result.nonEmpty)result.iteratorelseIterator.empty}else{buffer((name,score))Iterator.empty}}}四、方案 2分组 TopN4.1 ❌ 反模式groupByKeyrdd.groupByKey().mapValues(_.toList.sortBy(-_._2).take(3))// 每个 Key 的 Value 全部加载到单节点内存 → OOM4.2 ⚠ 中等方案reduceByKey 局部 TopNrdd.map{case(key,v)(key,List(v))}.reduceByKey((a,b)(ab).sortBy(-_._2).take(3))// Map 端先局部合并Combiner → Shuffle 数据量大幅减少// 但 RDD API 写起来较繁琐4.3 ✅ 推荐方案Spark SQL 窗口函数-- 每个班级取成绩 Top 3SELECTclass,name,scoreFROM(SELECT*,row_number()OVER(PARTITIONBYclassORDERBYscoreDESC)ASrnFROMstudents)tWHERErn3// DataFrame APIimportorg.apache.spark.sql.expressions.WindowvalwindowSpecWindow.partitionBy(class).orderBy($score.desc)df.withColumn(rn,row_number().over(windowSpec)).where($rn3).drop(rn)五、groupByKey vs reduceByKey 本质区别维度groupByKeyreduceByKeymapSideCombine❌ false✅ trueShuffle 数据量全部 KV 对Map 端预聚合后少量数据内存压力高全量加载低逐批合并适用场景需要对 Value 做非结合性操作可结合、可交换的聚合等价窗口函数collect_list()sum/count/max/min 自定义 UDAF// reduceByKey 的内部 Combiner 机制rdd.reduceByKey(__)// 源码等价于:// combineByKey(createCombiner, mergeValue, mergeCombiners)// Map 端先 mergeValue → mergeCombiners 分批合并六、实战每个班级成绩 Top 3// 方案 A: RDD reduceByKey较复杂rdd.map{case(cls,name,score)(cls,List((name,score)))}.reduceByKey((a,b)(ab).sortBy(-_._2).take(3)).flatMapValues(identity)// 方案 B: Spark SQL 窗口函数推荐valwWindow.partitionBy(class).orderBy($score.desc)df.withColumn(rn,row_number().over(w)).where($rn3)// 方案 C: rank/dense_rank处理并列// rank(): 1,1,3,4 (并列占位)// dense_rank(): 1,1,2,3 (并列不占位)// row_number(): 1,2,3,4 (纯粹行号)七、性能总结方案Shuffle 量内存安全推荐度groupByKey100%❌ 高危⛔ 禁止reduceByKey10-30%⚠ 可控可用窗口函数10-30%✅ 最优⭐ 首选金句groupByKey 是 Spark 新手的第一大坑——它把分组当成目的却忘了分组只是手段。真正的目的是在减少 Shuffle 的前提下拿到想要的结果。记住能用 reduceByKey 绝不用 groupByKey能用窗口函数就用窗口函数。作者starzy | AI Data Engineer / 大数据技术实践者博客blog.starzy.cn | GitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践