Spark二次排序与分组TopN性能优化实战:从原理到调优 📅 发布时间:2026/8/20 5:32:32 👁 浏览次数: 1. 项目概述从排序与取TopN说起在大数据处理领域Spark 以其卓越的内存计算能力和灵活的编程模型成为了处理海量数据的首选框架之一。然而随着数据量的激增和业务逻辑的复杂化一些看似基础的操作比如排序和分组取TopN却可能成为性能瓶颈的重灾区。我最近在优化一个实时报表系统时就遇到了一个典型的场景需要对每天数亿条的用户行为日志先按“用户ID”和“行为时间”进行二次排序再在每个用户分组内取出其最近发生的10条行为记录。这个需求听起来简单但初始的实现方案在集群上跑了近一个小时资源消耗巨大完全无法满足实时性要求。这促使我深入探究了Spark中二次排序和分组取TopN的底层实现机制并尝试了多种优化策略。经过几轮迭代最终将作业执行时间压缩到了十分钟以内资源消耗也大幅下降。今天我就把这次“踩坑”与“填坑”的全过程以及背后的核心原理和优化思路系统地梳理出来。无论你是刚接触Spark的新手还是已经有一定经验但想深入理解其内部机制的开发者相信这篇从实战中总结的干货都能给你带来启发。我们将从最朴素的实现开始一步步剖析性能问题所在并引入更高效的解决方案。2. 核心需求解析为什么二次排序和分组TopN是性能杀手在深入优化之前我们必须先理解这两个操作在Spark中为何消耗巨大。这不仅仅是写几行代码的问题而是关乎数据分布、Shuffle机制和内存管理的深刻理解。2.1 二次排序的本质与挑战二次排序顾名思义就是按照两个或以上的键进行排序并且主键和次键的排序规则可能不同。例如我们的需求是(userId, actionTime)先按userId升序再按actionTime降序。在Spark中最直观的做法是使用sortBy或orderBy。// 示例朴素但低效的二次排序 val rawRDD sc.textFile(“hdfs://path/to/user_actions”) val pairedRDD rawRDD.map(line { val fields line.split(“,”) (fields(0), fields(1).toLong) // (userId, actionTime) }) // 使用sortBy进行二次排序 val sortedRDD pairedRDD.sortBy(pair (pair._1, -pair._2))性能瓶颈分析全局ShufflesortBy或orderBy是一个全局排序操作。为了得到全局有序的结果Spark必须将所有分区的数据收集起来Shuffle Write然后在一个或少数几个Reducer中进行排序Shuffle Read。当数据量达到亿级时这个Shuffle过程会产生海量的磁盘I/O和网络传输成为最主要的性能瓶颈。全量数据移动即使你最终只需要每个分组的前N条这个操作也会移动和排序全部数据造成了巨大的资源浪费。内存压力负责最终排序的Reducer节点需要将属于它的所有数据加载到内存中进行排序极易导致Executor内存溢出OOM。注意sortBy默认是升序。对于次键降序我们通过-pair._2实现针对数值类型。对于非数值类型或更复杂的规则需要自定义排序器这本身也会带来一定的开销。2.2 分组取TopN的常见误区在得到排序后的RDD后新手常会这样取TopN// 示例低效的分组取TopN基于groupByKey val groupedRDD sortedRDD.groupByKey() val topNPerGroup groupedRDD.mapValues(iter iter.take(10).toList)性能瓶颈分析groupByKey的灾难groupByKey会将同一个键的所有值都拉取到同一个Executor上。如果一个“用户”的行为记录特别多比如一个活跃用户就会导致数据倾斜某个Task处理的数据量远大于其他Task拖慢整个作业甚至引起OOM。重复排序的浪费我们之前已经进行了全局排序但groupByKey本身不保证分组内值的顺序虽然因为上游RDD有序迭代器可能有序但这不是契约保证的。更重要的是全局排序的成本我们已经付出了而groupByKey又进行了一次全量的数据聚集。两阶段Shuffle整个流程经历了sortBy一次全量Shuffle和groupByKey又一次全量Shuffle相当于数据被大规模移动了两次效率极低。核心结论将“全局二次排序”和“基于groupByKey的分组TopN”组合使用是Spark作业中最常见的反模式之一。它产生了不必要的全量Shuffle和潜在的数据倾斜必须被优化。3. 优化策略一避免全局排序使用repartitionAndSortWithinPartitions我们的第一个优化方向是避免昂贵的全局排序。Spark提供了一个强大的转换操作repartitionAndSortWithinPartitions。它允许我们在进行数据重分区Shuffle的同时在每个分区内部进行排序。这比先repartition再sortBy高效得多因为排序和Shuffle的序列化/反序列化过程可以部分融合。3.1 自定义分区器与排序器要实现按userId分区并在分区内按actionTime降序排我们需要自定义分区器和排序规则。import org.apache.spark.Partitioner import org.apache.spark.rdd.RDD // 1. 自定义分区器按userId的哈希值分区 class UserIdPartitioner(override val numPartitions: Int) extends Partitioner { override def getPartition(key: Any): Int { // 假设key是 (userId: String, actionTime: Long) val userId key.asInstanceOf[(String, Long)]._1 math.abs(userId.hashCode) % numPartitions } } // 2. 隐式定义排序规则主键userId升序次键actionTime降序 implicit val sortOrdering: Ordering[(String, Long)] new Ordering[(String, Long)] { override def compare(x: (String, Long), y: (String, Long)): Int { val userIdCompare x._1.compareTo(y._1) if (userIdCompare ! 0) { userIdCompare } else { // 注意对Long类型直接比较会溢出。更安全的写法是 // java.lang.Long.compare(y._2, x._2) // 这里为了清晰使用数学比较 if (y._2 x._2) -1 else if (y._2 x._2) 1 else 0 } } } // 3. 应用优化 val optimizedSortedRDD: RDD[((String, Long), Null)] pairedRDD .map(pair (pair, null)) // 将数据包装成 (K, V) 形式V为null .repartitionAndSortWithinPartitions(new UserIdPartitioner(200)) // 假设设置200个分区 // 获取排序后的键 val resultKeysRDD optimizedSortedRDD.map(_._1)优化点解析单次Shuffle融合排序repartitionAndSortWithinPartitions将数据重分区和分区内排序合并为一个Stage减少了中间环节。数据在Map端按目标分区和排序键排序后溢出到磁盘在Reduce端读取时已经是分区有序的可以直接归并减少了Reduce端的排序压力。控制分区数通过numPartitions参数本例中为200我们可以根据集群资源和数据量合理设置避免分区过多或过少带来的任务调度开销或单分区数据过大问题。避免全量排序数据只在分区内有序全局是无序的。但对于后续“分组取TopN”来说只要同一个组的数据被分到了同一个分区并且在该分区内是排好序的就足够了。实操心得repartitionAndSortWithinPartitions要求RDD的键Key类型必须是可排序的Ordered。我们通过隐式Ordering对象来定义排序规则。自定义Ordering时比较逻辑要小心处理边界条件和溢出问题特别是对数值类型。3.2 结合mapPartitions高效取TopN现在我们得到了一个RDD其中相同userId的数据在同一个分区内并且按actionTime降序排列。接下来我们可以在每个分区内部高效地取出每个用户的TopN记录。这里使用mapPartitions是最高效的方式因为它允许我们在分区内维护一个状态如HashMap来累积TopN。val topNPerUserRDD optimizedSortedRDD.mapPartitions { iter // 使用一个可变的Map来存储每个userId的TopN列表 import scala.collection.mutable.{HashMap, ListBuffer} val topNMap new HashMap[String, ListBuffer[Long]]() iter.foreach { case ((userId, actionTime), _) val list topNMap.getOrElseUpdate(userId, ListBuffer[Long]()) if (list.size 10) { // 取Top 10 list actionTime // 如果是取TopN且需要保持顺序可以在这里插入排序。由于输入已按时间降序直接追加即可。 } // 如果list已满且actionTime比list中最小的还大对于降序即更晚的时间 // 则需要替换。但因为我们输入是降序的所以第一个遇到的就是最大的后续的都不会比已存的更大对于时间戳。 // 因此当list满后后续数据可以直接跳过。这是输入有序带来的关键优化 } // 将Map中的数据扁平化输出为 (userId, actionTime) 的迭代器 topNMap.iterator.flatMap { case (userId, timeList) timeList.map(time (userId, time)).iterator } }优化点解析零Shuffle整个TopN计算过程发生在每个分区内部通过mapPartitions完成没有引入任何额外的Shuffle操作。内存效率高每个Task处理一个分区只需要在内存中维护一个HashMap[String, ListBuffer]。假设一个分区包含1万个用户每个用户存10个Long内存开销大约为10000 * (key约50字节 10*8字节) ≈ 1.3MB完全可控。这远比groupByKey把整个分组数据拉取到内存要高效得多。利用输入有序性这是最关键的一点。由于分区内数据已经按actionTime降序排列当我们遍历迭代器时每个用户最先遇到的10条记录就是其Top10。一旦一个用户的TopN列表被填满后续所有属于该用户的记录都可以直接跳过大大减少了处理量。注意事项数据倾斜的缓解自定义分区器UserIdPartitioner基于hashCode分发数据能将不同用户均匀理想情况下分布到不同分区缓解了因单个用户数据过多导致的数据倾斜问题。但如果存在某些“热点用户”其数据量远超其他用户即使被哈希到一个分区该分区的Task仍可能负载过重。这是哈希分区的固有局限需要更高级的倾斜处理策略。分区数选择分区数 (numPartitions) 需要谨慎设置。通常建议设置为集群总核心数的2-3倍。过少会导致并行度不足且分区内数据量过大过多则会产生大量小任务增加调度开销。可以通过spark.default.parallelism参数来设定默认值并在作业中根据数据量调整。4. 优化策略二使用aggregateByKey与堆结构mapPartitions方案虽然高效但需要我们手动管理内存中的数据结构。另一种更函数式、且能灵活定义聚合逻辑的方法是使用aggregateByKey。结合可变集合如优先队列PriorityQueue我们可以在Shuffle的聚合阶段直接完成TopN的计算。4.1 利用aggregateByKey实现分区间聚合aggregateByKey需要三个参数一个初始值zeroValue一个用于分区内合并的函数seqOp和一个用于分区间合并的函数combOp。我们可以用PriorityQueue作为聚合的容器。import scala.collection.mutable.PriorityQueue val topN 10 // 定义一个函数用于向PriorityQueue中添加元素并保持其大小为TopN def addToHeap(heap: PriorityQueue[Long], value: Long): PriorityQueue[Long] { heap.enqueue(value) if (heap.size topN) { heap.dequeue() // 移除最小的元素对于最小堆队首是最小值 } heap } val topNByAggregateRDD pairedRDD.aggregateByKey(PriorityQueue[Long]())( // seqOp: 分区内合并将当前值合并到堆中 (heap, time) addToHeap(heap, time), // combOp: 分区间合并合并两个堆 (heap1, heap2) { heap2.foreach(time addToHeap(heap1, time)) heap1 } ).mapValues(heap heap.toList.sorted.reverse) // 将堆转为列表并排序降序原理与优化点Shuffle优化aggregateByKey会在Map端分区内先进行合并seqOp这称为“map-side combine”或“局部聚合”。这能显著减少Shuffle过程中需要传输的数据量。例如一个用户在一个分区内有100条记录经过seqOp后只会向Reduce端发送一个包含最多10条记录的PriorityQueue数据压缩率高达90%。灵活的聚合逻辑通过自定义seqOp和combOp我们可以实现复杂的聚合逻辑不仅仅是求和、计数。这里我们实现了基于堆的TopN选择。避免全量数据移动与groupByKey相比aggregateByKey传输的是聚合后的中间结果大小有限的堆而不是原始的全量数据。4.2 堆的选择与排序顺序上面的例子使用了Scala的PriorityQueue它默认是一个最大堆队首元素最大。但我们的需求是取最大的N个值即actionTime最大的TopN。在最大堆中队首是最大值。当我们维护一个大小为N的最大堆时队首是当前堆中的最大值但我们需要踢出的是“最小值”以保持堆里是最大的N个。因此我们需要一个最小堆Min Heap这样队首始终是当前堆中最小的元素即“第N大”的候选一旦有新元素比这个队首大就替换它。Scala的PriorityQueue可以通过传入一个Ordering来改变堆序。// 正确使用最小堆来维护最大的TopN implicit val minHeapOrder: Ordering[Long] Ordering[Long].reverse // 反转顺序使小的在前最小堆 val initialHeap PriorityQueue[Long]()(minHeapOrder) def addToMinHeap(heap: PriorityQueue[Long], value: Long): PriorityQueue[Long] { if (heap.size topN) { heap.enqueue(value) } else if (value heap.head) { // 如果新值比堆顶当前第N大大 heap.dequeue() // 移除堆顶最小的那个 heap.enqueue(value) // 加入新值 } heap }关键点最小堆维护最大TopN这是算法中的经典技巧。用最小堆保存当前遇到的最大的N个元素堆顶是这N个里最小的。新元素只需和堆顶比较即可决定是否替换。比较逻辑在addToMinHeap中value heap.head是关键。因为是最小堆heap.head是当前已保存的N个最大值中的最小值。如果新值比它还大说明新值有资格进入TopN。实操心得使用aggregateByKey时zeroValue这里是空的PriorityQueue会在每个分区内每个键的聚合开始时被创建。确保zeroValue是不可变的或者像我们这样在函数中返回一个新的集合。如果zeroValue是可变对象并在多个地方被修改会导致难以调试的数据错误。虽然我们的PriorityQueue是可变的但我们在每个seqOp/combOp中返回的是经过修改的新引用虽然物理对象相同但逻辑上是新的状态这种模式在Spark中是允许且常见的但需要非常小心。5. 优化策略三Spark SQL与窗口函数的降维打击如果你的数据源是结构化的例如Parquet、ORC、JDBC或者你已经将RDD转换成了DataFrame/Dataset那么Spark SQL提供的窗口函数Window Functions将是解决分组TopN问题的“终极武器”。它的声明式语法不仅简洁而且Spark的Catalyst优化器能够为其生成非常高效的执行计划。5.1 使用row_number()窗口函数假设我们有一个DataFramedf包含userId和actionTime两列。import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ val windowSpec Window.partitionBy(“userId”).orderBy(col(“actionTime”).desc) val resultDF df .withColumn(“rn”, row_number().over(windowSpec)) .where(col(“rn”) 10) .drop(“rn”)短短四行代码发生了什么Window.partitionBy(“userId”)定义了窗口数据将按userId分组。.orderBy(col(“actionTime”).desc)在每组内按actionTime降序排列。row_number()为每组内排序后的每一行分配一个唯一的序号1, 2, 3…。where(col(“rn”) 10)过滤出序号小于等于10的行即每个用户的Top10。drop(“rn”)丢弃临时序号列。5.2 Spark SQL执行计划与优化我们可以通过resultDF.explain()查看物理执行计划。你会看到类似如下的内容简化后 Physical Plan *(4) Project [userId#0, actionTime#1L] - *(4) Filter (isnotnull(rn#8) (rn#8 10)) - Window [row_number() windowspecdefinition(userId#0, actionTime#1L DESC NULLS LAST, specifiedwindowframe(RowFrame, unboundedpreceding$(), currentrow$())) AS rn#8], [userId#0], [actionTime#1L DESC NULLS LAST] - *(3) Sort [userId#0 ASC NULLS FIRST, actionTime#1L DESC NULLS LAST], false, 0 - Exchange hashpartitioning(userId#0, 200) - *(2) ...优化器带来的优势智能的Shuffle计划中只有一个Exchange即Shuffle按userId进行哈希分区。这与我们手动使用repartitionAndSortWithinPartitions的思路一致。排序下推Sort操作发生在Shuffle之后在每个分区内进行。由于窗口函数要求数据在分区内有序Spark会自动安排这个排序并且可能与其他操作如Shuffle的序列化进行优化。管道化执行整个计算过程可以被Spark的Whole-Stage Code Generation全阶段代码生成优化编译成高效的Java字节码避免虚拟函数调用等开销性能通常优于手写的RDD算子。与RDD方案的对比代码简洁性SQL/DataFrame API的代码极其简洁意图清晰。性能对于大多数场景经过Catalyst优化器和Tungsten执行引擎优化后的SQL方案性能与手动优化的RDD方案持平或更优尤其是在处理结构化数据时。灵活性对于极其复杂、非标准化的聚合逻辑RDD的aggregateByKey可能更有优势。但对于标准的排序、TopN、排名等操作窗口函数是不二之选。注意事项窗口函数在计算row_number、rank、dense_rank时需要将整个窗口的数据物化到内存中吗不一定。对于row_number由于只需要顺序编号Spark可以实现流式处理内存消耗只与分组键的数量有关需要跟踪当前序号而与组内数据量关系不大。但对于rank和dense_rank在遇到相同值时需要特殊处理内存占用可能会稍高。无论如何它都比groupByKey全量拉取数据要高效得多。6. 高级优化与调优实战掌握了核心策略后我们还需要关注一些高级调优技巧以应对极端数据规模或特殊场景。6.1 应对数据倾斜两阶段聚合与盐析技术当某个userId的数据量异常庞大时例如一个测试用户或爬虫用户即使用哈希分区这个用户的所有数据也会被分到同一个分区导致该分区Task运行缓慢。解决方案局部聚合全局聚合两阶段聚合第一阶段加盐局部聚合给每个键附加一个随机前缀盐将原本的一个热点键分散成多个临时键。val saltedRDD pairedRDD.map { case (userId, time) val salt (math.random * 10).toInt // 生成0-9的随机盐 ((salt, userId), time) // 键变为 (salt, userId) } // 对 (salt, userId) 进行聚合取TopN。此时原热点键的数据被分散到10个盐值上。 val partialTopN saltedRDD.aggregateByKey(...) // 使用前面定义的aggregateByKey逻辑但键是(salt, userId)第二阶段去盐全局聚合去掉盐前缀对同一userId来自不同盐值的中间结果进行最终聚合。val finalTopN partialTopN .map { case ((salt, userId), heap) (userId, heap) } // 去掉盐 .reduceByKey((heap1, heap2) mergeHeaps(heap1, heap2)) // 合并各盐值的结果 .mapValues(heap heap.toList.sorted.reverse)原理通过加盐我们将一个大的数据块打散到多个Task中并行进行局部TopN计算。虽然增加了额外的Shuffle阶段reduceByKey但有效避免了单个Task长时间运行总体耗时往往更低。6.2 内存与GC优化无论是mapPartitions中的HashMap还是aggregateByKey中的PriorityQueue都涉及JVM堆内对象的管理。在数据量极大时GC可能成为瓶颈。使用原始数据类型如果可能尽量使用Long、Double等原始类型避免使用包装类型java.lang.Long和复杂的对象以减少内存占用和GC压力。在Scala中使用specialized注解或直接使用数组存储long可能更高效但代码更复杂。调整Spark内存配置spark.executor.memory增加Executor总内存。spark.memory.fraction/spark.memory.storageFraction调整用于执行和存储的内存比例。对于这类Shuffle和聚合密集的作业可以适当提高spark.memory.fraction默认0.6。spark.sql.windowExec.buffer.spill.threshold当使用窗口函数时这个参数控制窗口聚合缓冲区的溢出阈值。如果遇到OOM可以调低此值默认4096让Spark更早地将中间结果溢出到磁盘。序列化优化如果自定义的键或值类非常复杂考虑使用Kryo序列化spark.serializer-org.apache.spark.serializer.KryoSerializer并注册自定义类这能减少Shuffle时的序列化体积和耗时。6.3 分区数动态评估与自适应查询执行在Spark 3.x中自适应查询执行AQE功能可以动态优化运行时执行计划。spark.sql.adaptive.enabled设置为true以启用AQE。spark.sql.adaptive.coalescePartitions.enabledAQE可能会在Shuffle后合并过小的分区避免大量小任务。spark.sql.adaptive.skewJoin.enabledAQE可以自动检测并处理Shuffle Join中的数据倾斜。对于RDD API虽然不能直接享受AQE对SQL的优化但我们可以从AQE的思路中获得启发监控Stage的Task执行时间如果发现最大最小时间差很大很可能存在数据倾斜就需要手动介入采用前面提到的盐析等技术。7. 方案对比与选型指南至此我们介绍了三种核心优化策略及其变种。下面用一个表格来总结对比帮助你在实际场景中做出选择。特性/方案朴素方案 (sortBy groupByKey)优化方案1 (repartitionAndSort mapPartitions)优化方案2 (aggregateByKey Heap)优化方案3 (Spark SQL Window)核心思想全局排序后全量分组分区内排序分区内局部TopNMap端Combiner局部聚合减少Shuffle数据量声明式编程由Catalyst优化器生成计划Shuffle次数2次 (全量排序 全量分组)1次 (重分区并排序)1次 (按Key聚合)1次 (按分区键Shuffle)数据移动量极大 (全量数据移动两次)中等 (全量数据移动一次)较小(仅传输聚合后的堆)中等 (全量数据移动一次)内存压力极大 (Reducer端全量数据)较小 (Task内维护每个Key的TopN列表)小 (Task内维护每个Key的堆)取决于窗口函数实现通常较小代码复杂度低中 (需自定义分区器/排序器)中 (需理解聚合函数和堆)极低灵活性低高 (可完全控制分区和计算逻辑)高 (可自定义复杂聚合逻辑)中 (受限于SQL表达式和UDF)适用场景不推荐仅用于学习反面案例需要精细控制分区、排序逻辑的复杂RDD作业Key数量多单个Key的Value数量大需要高效聚合的RDD作业首选数据源为结构化数据逻辑为标准排序/排名抗数据倾斜差依赖分区器可通过盐析增强依赖分区器Map端Combine可缓解AQE可自动优化也可手动盐析选型建议如果你的数据已经是DataFrame/DataSet或者可以轻松转换毫不犹豫地选择Spark SQL窗口函数方案方案3。这是开发效率、运行效率和可维护性的最佳平衡。如果你需要处理非结构化数据或者有非常复杂、非标准的聚合逻辑无法用SQL表达那么RDD API更合适。如果分组键数量可控且TopN的N较小mapPartitions方案方案1非常高效。如果分组键数量巨大或者单个键对应的值非常多aggregateByKey方案方案2利用Map端Combine的优势更大能极大减少Shuffle数据量。永远避免使用sortBygroupByKey的组合。8. 实战问题排查与性能调优记录在实际操作中我遇到了几个颇具代表性的问题这里分享排查过程和解决方法。问题一作业卡在最后一个Stage某个Task运行时间远超其他Task。现象使用aggregateByKey方案大部分Task在1分钟内完成但总有1-2个Task运行超过20分钟。排查查看Spark UI的Stages页面找到慢Task所在的Stage。查看该Stage的Event Timeline发现慢Task的“Shuffle Read Size”远大于其他Task。结论数据倾斜。某个或某几个userId的数据量极大导致对应的分区数据量暴增。解决采用6.1中所述的盐析技术。为Key添加随机前缀打散热点Key。调整盐的数量例如0-99需要在增加Shuffle开销和平衡负载之间取得平衡。通过实验选择使最慢Task时间缩短至平均Task时间2倍以内的盐值数量。问题二使用窗口函数时出现Executor OOM。现象java.lang.OutOfMemoryError: Java heap space错误栈指向org.apache.spark.sql.catalyst.expressions.GeneratedClass。排查检查输入数据发现存在大量userId为null的记录。窗口函数partitionBy(“userId”)会将所有null值分到同一个分区造成严重倾斜。另外窗口的orderBy子句如果存在大量重复值且使用了rank()或dense_rank()可能需要维护更多状态。解决数据清洗在应用窗口函数前过滤掉userId为null或无效值的记录。df.filter(col(“userId”).isNotNull)。调整配置增加Executor内存 (spark.executor.memory)。同时调低窗口函数缓冲区的溢出阈值spark.sql.windowExec.buffer.spill.threshold例如从4096设为1024让Spark更激进地将中间状态溢出到磁盘用时间换空间。考虑替代方案如果倾斜无法消除且row_number()仍OOM可以回退到RDD的aggregateByKey方案并对其应用盐析技术以获得更精细的内存控制。问题三repartitionAndSortWithinPartitions后数据似乎没按预期排序。现象在mapPartitions中打印数据发现同一个userId的数据不是严格按照actionTime降序排列。排查检查自定义的Ordering比较逻辑发现用于比较Long的代码有误在差值超过Int范围时可能返回错误结果。检查分区器getPartition逻辑确保它只依赖于userId而不是整个元组(userId, actionTime)。如果依赖了整个元组那么(userId, time1)和(userId, time2)可能被哈希到不同分区导致同一个用户的数据分散开分区内排序自然无法保证该用户数据的全局有序但分区内跨用户是有序的。解决使用java.lang.Long.compare(x, y)来安全地比较Long值。确保分区器仅基于分组键userId计算分区号。这是repartitionAndSortWithinPartitions能正确工作的前提相同分组键的数据必须进入同一个分区。经过这一系列的剖析、优化和实践最初需要运行近一小时的作业最终稳定在10分钟以内完成。这个优化过程深刻地告诉我在Spark中理解数据流动Shuffle和内存管理的成本远比写出能跑通的代码更重要。选择正确的算子往往能带来数量级的性能提升。下次当你面对排序和TopN需求时不妨先停下来想想是否真的需要全局排序数据是否倾斜能否在Shuffle前就进行压缩用这些问题指引你的优化方向你的Spark作业一定会跑得更快、更稳。