1. 一次慢任务复盘为什么40分钟的作业最后只优化到3分钟先说个真实经历。之前接手一个离线数据处理任务凌晨跑批逻辑也不复杂就是把用户行为日志按维度聚合后落表。但作业运行时间越来越长从最初的20分钟涨到40多分钟集群资源也没加过数据量增幅也不算夸张。排查时发现罪魁祸首就是一堆写法随意的Spark算子一个groupByKey后跟map做求和一个collectAsMap把大表拉到Driver端还有两处filter之后又repartition再filter。说白了不是RDD本身慢是算子用错了地方。这也是我写这篇文章的初衷。SparkCore里的算子如果只是翻文档看签名你只会觉得哦这个API能做什么但放在真实作业里同样的功能用不同算子实现性能和稳定性天差地别。RDD是Spark最基本的数据抽象而算子就是操作这个抽象的刀。刀用得好不好直接决定作业跑得顺不顺。这篇文章想聊的东西很实在RDD算子到底分成哪几类、每一类的执行逻辑是什么、哪些算子会拖垮性能、遇到常见场景该怎么选算子。适合刚入门的Spark开发者也适合写过一段时间但总是被作业性能困扰的同学。看完你至少能回答一个问题你的作业慢到底是慢在哪个算子、哪一步shuffle。2. 把算子先分家Transformation、Action与控制算子2.1 两大门派惰性的Transformations与贪婪的ActionsRDD算子第一个要搞明白的分水岭就是Transformation和Action的区别。这直接决定了你对Spark执行模型的理解。Transformation是对RDD做变换比如map、filter、flatMap、distinct、reduceByKey。这类算子的特点是不触发实际计算只是构建一个RDD的依赖关系链。换句话说你对一个RDD调用mapSpark并不会立刻遍历数据而是记下我要对每个元素做这个变换这个血缘信息。Action则是真正触发计算的算子比如count、collect、reduce、saveAsTextFile、foreach。只有遇到ActionSpark才会把之前积累的所有Transformation串成一个DAG划分成Stage然后提交执行。我见过不少新手刚开始写Spark作业每一个步骤都习惯加一个count()来看数据条数。结果就是一个本来能一趟跑完的作业被强行切成了好多段每个count都会触发一次完整的DAG计算。这在调试小数据时无所谓一旦跑到亿级数据每次count都是一次全量扫描作业不慢才怪。2.2 容易被忽略的第三类控制算子除了Transformation和ActionSparkCore里还有一类算子经常被忽略——控制算子。主要是cache、persist、checkpoint。它们不改变RDD的数据内容而是控制RDD的存储方式和生命周期。cache和persist的作用是把中间结果持久化到内存或磁盘避免同一个RDD在多个Action中被重复计算。checkpoint则会把RDD的数据真正保存到可靠存储通常是HDFS并且切断血缘链。这三类算子在长作业、复杂DAG中作用极大后面第6章我会专门展开讲。2.3 第一课小结先记住这个判断框架看到一个算子第一反应不是这个算子能干什么而是这个算子是Transformation还是Action。这个判断决定了它会不会立刻触发计算也决定了它出现在DAG中的位置。后续所有关于慢不慢的分析都从这个分类开始。3. Transformation系列拆解惰性背后的执行细节3.1 单元素变换map与flatMap的边界怎么划map和flatMap可能是你用得最多的两个算子。它们看起来很像都是逐元素操作但返回结构不同map是一对一每个输入元素产生一个输出元素flatMap是一对多每个输入产生零到多个输出并把结果扁平化。举一个真实的例子。处理用户行为日志时日志的JSON字段里可能有一个数组比如浏览的商品列表你想把每条日志展开成一行为一个商品。这时候如果写map得到的是一个包含List的RDD你还得再flatMap一次才能拆开。如果直接用flatMap一步就到位了。这类操作在ETL清洗里非常常见。// 错误示范先用map再手动拍平多出一次算子调用 val itemsRDD logsRDD .map(log extractItems(log)) .flatMap(list list) // 推荐写法flatMap一步完成 val itemsRDD logsRDD.flatMap(log extractItems(log))虽然从血缘上看这两种写法最终都会被Spark优化但多一步算子就是多一层嵌套函数调用代码可读性也差。能用flatMap表达的语义别硬拆成两步。3.2 聚合家族reduceByKey与groupByKey的真实差距这是经典到不能再经典的话题但依然值得展开讲因为它是区分会用算子和懂Spark的分界线。先看现象。下面两行代码输出结果一样性能差一个数量级// 方式AreduceByKey 先本地聚合 val wordCountA wordsRDD.map(w (w, 1)).reduceByKey(_ _) // 方式BgroupByKey 后做map聚合 val wordCountB wordsRDD.map(w (w, 1)).groupByKey().map { case (k, vs) (k, vs.sum) }为什么差别这么大关键在于shuffle的数据量。reduceByKey是分区内先合并每个分区内同一个key的多个记录先加成一两个中间结果然后再做跨节点的shuffle。groupByKey则不做任何本地合并直接把所有相同的key的values原封不动地通过网络传到目标节点再由下游算子做聚合。数据量越大groupByKey要传输的数据就越多I/O和网络开销就越大。我跑过一个几百GB的日志聚合任务。把groupByKey().map(sum)换成reduceByKey后shuffle数据量从大约120GB降到了18GB作业时间缩短了将近70%。这不是微调是质变。有一个注意点并不是说groupByKey一无是处。如果你不仅要聚合值还要对同一个key的整个value集合做复杂的非交换非结合操作比如求中位数、取Top N中依赖全部元素的操作这时候groupByKey依然有它的价值。但绝大多数按key求和、计数、取最大这类场景都应该优先考虑reduceByKey或aggregateByKey。顺带说一句foldByKey和aggregateByKey也是聚合家族成员。aggregateByKey是更通用的存在因为它在分区内和跨分区合并时可以使用不同的函数这在进行先求局部平均值再合并全局这种操作时很实用。3.3 过滤算子filter的位置决定了扫描的数据量filter看起来简单但它的位置安排对性能影响很大。原则是越早过滤越好。因为Spark的计算模型是每个Transformation都会对全部分区数据做一次函数调用你在前面过滤掉的数据后面所有算子都不用再处理了。实际作业里我经常看到这样的写法val result sourceRDD .map(parseLog) .flatMap(extractTags) .map(enrichData) .filter(isValid)如果isValid判断可以提前到parseLog之后就应该提前。比如日志里大量字段缺失的记录直接在最开始就filter掉后面的flatMap和enrichData能少处理几十个百分点无效数据。还有一个小技巧能用filter的地方别总想着repartition。有些人觉得数据倾斜就repartition打散结果shuffle一次数据量大的时候反而更慢。先filter到小数据量再考虑要不要调整分区这个顺序别反了。3.4 去重与集合操作distinct、union、intersection的性能提醒distinct底层是通过reduceByKey实现的本质上是一次shuffle去重。数据量大的时候distinct的成本不低因为要跨节点识别重复元素。union是最廉价的集合算子因为它只是把两个RDD的分区拼在一起不涉及任何计算和数据混洗。但要注意union后RDD的分区数是两个RDD分区数之和如果两边分区数都不小下游任务数量会翻倍影响并行度设置。intersection和subtract都会触发shuffle因为它们需要跨分区比较元素是否存在。能用join替代的场景通常join更高效——因为join的shuffle是按key分组的而集合操作是基于元素完整值的比较数据量大时开销更大。4. Action系列触发机制DAG是怎么跑起来的4.1 DAG调度从Action反推Stage链Action的关键作用不是计算结果本身而是点燃整个DAG。当你对最终RDD调用一个Action时Spark会沿着血缘反推找出所有需要计算的父RDD将这些RDD按照宽窄依赖切分成Stage。这里必须理解宽依赖和窄依赖的概念因为Stage划分就是依据这个。窄依赖是指父RDD的每个分区最多被子RDD的一个分区使用比如map、filter、union。宽依赖是指父RDD的每个分区会被子RDD的多个分区使用典型代表是groupByKey、reduceByKey、repartition、join。宽依赖的出现意味着shuffle而shuffle就意味着Stage的边界。每一个宽依赖都是一个断点把DAG切断成独立的Stage。窄依赖的算子可以流水线执行多个算子合并在一个Stage里跑完宽依赖必须等上游Stage所有任务执行完才能开始下游Stage。4.2 常见Action算子的真实代价collect是最容易踩坑的Action。它会把所有分区的结果拉到Driver端然后组成一个数组。数据量大时Driver内存直接爆掉。我见过一个同学对全量日志执行collect然后自己遍历统计结果Driver OOM作业失败。正常做法是先用reduce或aggregate在Executor端完成统计只把很小的结果返回到Driver。take(n)用起来也要小心它的底层会先跑一个分区看够不够n个不够再跑更多分区如果数据分布极度不均匀它可能被迫跑完全部分区。saveAsTextFile这类输出型Action是唯一值得把全量结果写出去的操作但要注意输出路径目录数量。saveAsTextFile的输出文件数等于分区数如果你有几百个分区就会产生几百个小文件。小文件问题不只在HDFS下游数仓读取也会受影响。遇到这种情况在save之前先coalesce或repartition成合理分区数是很有必要的。4.3 foreach与foreachPartition连接外部系统的正确姿势foreach和foreachPartition是处理每个分区数据要写到外部系统的常用Action。直接foreach每个元素意味着每条记录都跟外部系统建立一次连接性能极差。正确做法是用foreachPartition在每个分区上建立一次连接批量写入。rdd.foreachPartition { partition val conn createConnection() partition.foreach { record conn.write(record) } conn.close() }这个模式在写Elasticsearch、MySQL、Kafka时都适用。单个分区的连接复用能把外部系统压力从每条记录一次连接降到每个分区一次连接。数据记录上千万时这个优化能减少几个数量级的连接数。5. 宽依赖与窄依赖作业性能的底层密码5.1 为什么说算子数量多不等于作业慢网上常有大量使用算子对硬件性能的挑战这种讨论但Spark里真正的性能挑战从来不是算子数量多而是算子链中是否出现了不必要的宽依赖。窄依赖的算子再多Spark都能把它们合并进同一个Stage流水线执行几乎可以看作一次遍历。真正打断流水线、引发网络传输和磁盘I/O的是每个宽依赖点。你可以把窄依赖想象成在一条传送带上加工零件上游没加工完下游马上接着处理中间不需要停顿。而宽依赖就像每个零件都先要送到中央仓库按照某种规则重新分拣再装上另一条传送带。分拣过程shuffle是整条流水线的瓶颈。所以判断一个作业优劣不要数算子个数要数宽依赖点数。很多作业慢就是宽依赖太多Stage链过长每个Stage之间的shuffle都成了拖累。5.2 shuffle机制里的几笔账数据量、IO次数与序列化开销shuffle的代价可以从三个维度理解。第一是数据量。shuffle要把key相同的数据从不同节点汇聚到同一节点网络传输的数据量取决于shuffle写的记录数。reduceByKey在map端做合并直接压缩了shuffle数据groupByKey完全不合并shuffle数据就是所有记录原样的总和。第二是IO次数。shuffle过程中每个Map任务都要写本地磁盘文件Reduce任务要读取这些文件。大量的小文件读写会拖垮磁盘性能。这也是为什么shuffle后要合理设置分区数分区太少单个任务数据量过大分区太多小文件爆炸。第三是序列化开销。shuffle传输的数据需要序列化和反序列化。使用Kryo序列化器能显著减少序列化后的字节数对shuffle性能是实打实的提升。在spark-defaults.conf里设置spark.serializerorg.apache.spark.serializer.KryoSerializer并注册需要用到的类属于低成本高收益的优化。5.3 大数据量Task的硬件压力Executor资源配置与分区数的平衡算子链对硬件资源的压力最终都体现在任务执行上。一个常见的误区是盲目加大spark.executor.memory以为内存大了就万事大吉。实际上如果分区数没变内存加大只是让每个任务在更宽松的环境下跑GC时间反而可能增加。合理思路是配合分区数调整。一个实用参考每个Executor的core数乘以Executor数量得到总并发slot合理的分区数建议是总并发slot的2到4倍。分区太少CPU利用不满分区太多任务调度和shuffle开销增大。遇到过实际案例某个作业分区数设为2Executor有20个core等于只用了10%的并行度。把分区数改成200后运行时间从30分钟降到6分钟。另外大任务里OOM的常见原因是groupByKey或aggregateByKey后单个key的数据量太大导致某个Reducer内存爆掉。此时应该从业务上考虑是否能把key拆分加盐salting或者退而求其次用countApproxDistinctByKey这类近似算子替代精确去重。对硬件而言这些手段的本质是把内存压力分摊到更多任务上。6. 控制算子cache、persist与checkpoint的正确打开方式6.1 同一个RDD被多个Action复用cache的典型场景默认情况下Spark对RDD的每一次Action都会重新计算一遍血缘链上的所有Transformation。如果一个RDD会被后续多个Action使用比如先count看一眼再saveAsTextFile输出那这个RDD就被算了两次。避免重复计算的办法就是cache。val normalizedRDD rawRDD .filter(validLog) .map(parseLog) .cache() val totalCount normalizedRDD.count() normalizedRDD.saveAsTextFile(hdfs:///output)缓存级别方面MEMORY_ONLY是默认选项适合数据量不大且反复读的场景。如果数据量较大内存放不下MEMORY_AND_DISK会溢写到磁盘代价是再次读取时多一次磁盘I/O。DISK_ONLY适合数据算出来代价很高、但内存又放不下的情况。这里有个判断标准缓存前的数据有多大预计被调用几次如果这个RDD只被用一次cache就是纯浪费。我看到很多新手习惯性地每个RDD都cache结果缓存占用大量内存反而让GC压力变大得不偿失。6.2 checkpoint剖断血缘链的救命手段checkpoint和cache不一样。checkpoint是把RDD的计算结果写到可靠存储HDFS并且切断之前的血缘关系。同一个RDD在checkpoint之后Spark认为它的数据已经固化了不需要再从头计算。它和cache最大的区别在于血缘。cache保留了完整血缘如果缓存丢失Spark还能根据血缘重新计算。checkpoint则直接斩断了血缘链缓存丢失也无法追溯上游了。那为什么要切断血缘对于非常长的DAG血缘链本身会占用Driver内存而且一旦某个很深的依赖链出错重算成本极高。checkpoint之后Spark可以直接从那一个点开始算血统一目了然。实际使用checkpoint要注意必须在checkpoint之前对RDD执行一个Action否则checkpoint不会真正触发。这个坑我自己踩过先rdd.checkpoint()然后又对rdd做各种Transformation最后直接save发现checkpoint根本没执行因为checkpoint是lazy的。正确做法是rdd.checkpoint(); rdd.count()强制它落盘然后再继续后续处理。6.3 缓存策略选型的一张决策表场景推荐策略原因同一RDD被多个Action使用且数据量适中MEMORY_ONLY读取快全内存数据量略大内存可能不够MEMORY_AND_DISK溢写磁盘兜底不丢失数据数据量巨大且重算代价高checkpoint到HDFS免疫故障砍断血缘算子链极长且涉及多次失败重试checkpoint在中间节点缩短重算路径RDD只用一次不缓存缓存本身有序列化和内存开销7. 实战选型经验从业务意图反推算子组合7.1 一个典型ETL任务的算子编排把前面讲的内容串起来我用一个典型的ETL任务做完整演示。需求读取一小时的行为日志清洗无效记录提取用户ID和商品ID按用户ID聚合出浏览过的商品列表最后输出到HDFS。算子编排可以是这样的val rawRDD sc.textFile(hdfs:///logs/hour2024-01-01-00) // 第一层尽量早过滤减少后续数据量 .filter(line line.contains(\user_id\) line.contains(\item_id\)) // 第二层解析JSON提取关键字段 .map(line parseLog(line)) // 第三层把有效记录保留过滤解析失败的数据 .filter(parsed parsed.isDefined) .map(parsed (parsed.userId, parsed.itemId)) // 第四层按用户聚合商品列表 .groupByKey() // 第五层转换格式准备输出 .map { case (userId, items) s$userId\t${items.mkString(,)} } // 输出前合并分区避免小文件 .coalesce(20) // Action触发 .saveAsTextFile(hdfs:///output/user_items)这个例子里有个可以优化的点groupByKey把每个用户的所有商品ID都放到内存里如果单个用户浏览记录特别多很容易OOM。如果业务上不需要完整的商品列表顺序只是需要聚合结果可以用distinct().reduceByKey(_ _)之类的方式但更稳妥的做法是先map成字符串格式再reduceByKey合并这样合并逻辑在shuffle前先执行了一部分内存压力会小很多。7.2 算子选择决策要点结合日常高频场景我整理出一个简单的决策要点清单先想清楚目标是变换数据还是触发计算Transformations vs Actions。聚合优先考虑reduceByKey、aggregateByKey而不是groupByKey。filter和distinct要放在DAG尽可能靠前的位置。需要把RDD写到外部系统时优先foreachPartition批量写而不是foreach逐条写。同一个中间结果被多个Action使用时再考虑cache或persist。遇到极长血缘链或频繁失败的任务在关键节点加checkpoint。输出前根据目标路径合理合并分区数控制小文件数量。7.3 一个被算子数量误导的真实案例之前讨论网络热词大量使用算子对硬件性能的挑战时团队里有人觉得是算子链太长导致CPU压力大。我们做了一次实验同一个数据源一个任务写了12个窄依赖算子一个任务写了5个算子但其中含两个groupByKey。结果第一批作业跑了4分钟第二个跑了23分钟。这足以说明硬件压力主要来自shuffle带来的网络和磁盘I/O窄依赖算子再多也只是串行执行在同一个Stage里的函数调用反倒是宽依赖算子一个就足够拉爆整个作业。所以以后在代码评审里看到有人担心算子太多影响性能我一般会反问你的DAG里有几个shuffle点这才是真正要盯的地方。8. 写在最后的一点个人体会单独把SparkCore算子拎出来讲是因为我发现很多人的问题不在不会写Spark而在不知道每个算子背后到底发生了什么。map、filter写多了形成了肌肉记忆但一旦遇到groupByKey和reduceByKey这种名字相近、功能相近、底层天差地别的算子就容易犯错。从我自己的经验看掌握RDD算子的核心不是背熟API列表而是建立两个思维习惯第一每个算子进DAG时是惰性挂起还是立即触发第二每个算子会不会引入宽依赖、产生shuffle。把这两个问题想清楚一半的性能问题和一半的OOM问题都能在设计阶段避免。这篇文章里讲的reduceByKey对比groupByKey、foreachPartition替代foreach、cache和checkpoint的取舍、分区数跟Executor资源的配合都是我在真实作业里反复试过、反复踩过之后沉淀下来的经验。你要是也正在为Spark作业的运行时间和集群稳定性头疼不妨先不看复杂的参数调优文档回来把算子层面的选择重新捋一遍。往往改动不大收益却很直接。
阅读完成 · 觉得有帮助?