做大数据的人,只要经历过MapReduce那个年代,大概都忘不了被“跑一批任务等到天亮”支配的感觉。那时候处理几百GB的数据,凑一个稳定跑完的作业都算技术活,更别说什么实时性、交互式查询。后来Spark出现,内存计算的概念打出来,很多团队像抓住救命稻草一样从MR迁到Spark,但也有人换了Spark以后发现作业照样慢、照样OOM,跑起来甚至不如调优过的MapReduce稳。
这篇文章不聊空泛的架构演进,全部围绕实际调优经验来写。从MapReduce到Spark,本质不是换一个框架,而是换一套对计算资源、数据Shuffle、内存模型的理解方式。你会看到为什么MapReduce慢、慢在哪个环节;Spark快、快在哪里但坑又藏在哪里;集群搭好之后哪些参数是真正影响性能的;Shuffle怎么做才能不拖后腿;Spark OOM怎么一步步排查;还有Spark SQL处理外部数据源(例如集成Redis、接入达梦数据库)时那几个容易忽略的IO瓶颈。内容适合刚入门大数据、准备给现有批处理环境做升级的工程师,也适合天天调Spark作业但始终觉得没摸到门道的朋友。
1. 为什么性能调优绕不开“从MapReduce到Spark”这段路
1.1 先明确要解决的真实问题
先定义一个问题:我们说的性能调优,到底在调什么?从MapReduce那个时代走到Spark的时代,表面上是换框架、换API,本质上是换了一套计算范式。MapReduce把每个计算都拆成Map和Reduce两个阶段,数据每一轮都要写入磁盘,下一轮再从磁盘拉回来。这个模型的好处是简单、稳定、容错容易做,坏处也非常直观:计算过程中全是磁盘IO和网络IO,时间都花在搬运上。
举个我实际接触过的案例。早年间做过一个用户行为日志的分析任务,输入是HDFS上大约300GB的文本日志,要做一次简单的词频统计和会话切分预处理。如果用MapReduce,Map阶段读日志吐中间结果,Sort和Shuffle阶段把数据按照Key排序分发,Reduce阶段再拉数据聚合成最终结果。整个过程走一遍,在大概10台物理机的集群上跑,耗时大约是以小时为单位计算的。后来同样的逻辑迁到Spark RDD来写,Stage内部能内存计算的部分全部走内存,跑完一次作业的时间缩短到分钟级别。
这个对比其实并不在于谁的算法更高级,而是Shuffle和磁盘落地的次数相差太多。MapReduce里的每个Job几乎强制要求中间结果落盘,因为它的容错模型就是“每个Stage重新执行”。Spark通过RDD的血统机制做容错,Stage内部能缓存的中间结果直接保留在内存或按需计算,不再强制落盘。这就是性能差距的一个根本来源。
1.2 性能瓶颈到底在哪:重新理解MapReduce的三个“太重”
之前有人问过我一个问题:MapReduce是不是被时代淘汰了?我的看法是,这个模型没有淘汰,但实现思路确实过重。重在哪里,拆开看有三点:
第一,中间结果落盘太重。Map结束后,MapOutput要写到本地磁盘,Reduce任务启动后再从各个Map任务节点拉取数据,这个过程既是磁盘IO又是网络IO。数据量一大,整个作业的时间几乎都消耗在落盘和拉取上。
第二,调度模型太重。MapReduce的每个Task是JVM级别的进程,意味着一个Job包含几千个Task,就要反复启动几千个JVM。JVM启动本身的开销、序列化、反序列化、上下文切换,这些和业务逻辑无关的额外消耗占比不小。Spark引入Executor常驻内存的模型,Task以线程方式在同一进程内调度,启动开销大幅下降。
第三,编程范式太重。MapReduce想表达一个多步骤的复杂计算,就得把逻辑拆成多个Job串联,每个Job之间通过HDFS传递中间数据,这种串联模式让复杂分析逻辑代价极高。而Spark用RDD或DataFrame天然支持在单个作业里构建多Stage的DAG,Stage内部尽可能走内存管道化计算。从需求和成本的角度讲,这个演进是必然的。
2. 环境搭建:Spark集群部署的参数取舍
2.1 版本选型与基础配置
Spark集群搭建这件事,听起来简单,真正影响性能的细节经常藏在安装之外。很多教程写spark安装步骤,无非是下载二进制包、配置spark-env.sh、slaves文件、启动集群,但这几步只能说让集群能跑,离“能跑得快”还差得远。
版本选型方面,我建议直接用当前主线稳定版本,比如Spark 3.4.x或3.5.x系列,原因是3.0之后的Spark引入了Adaptive Query Execution(AQE)等自动优化能力。AQE可以在运行时根据实际Shuffle输出统计信息重新优化执行计划,比如自动处理数据倾斜时拆分Reduce分区、自动合并过小的Shuffle分区,这些功能对刚入手调优的人来说非常友好,能减少不少手工调参的压力。
配套生态需要提前想清楚。如果你主要跑批处理,用YARN作为资源调度会比Standalone模式扩展性更好,别图省事长期跑Standalone。如果只是本机学习或做小规模实训,那Standalone加本地目录没问题。所谓spark集群搭建的“集群”二字,起码要覆盖NameNode、DataNode、ResourceManager、NodeManager、Spark的Executor运行节点这些角色,否则你在物理环境会踩到很多资源协商上的坑。
以我曾经搭过的一个测试环境为例,配置大致是这样:
| 节点角色 | 配置 | 用途 |
|---|---|---|
| 主节点 | 16核 / 64GB内存 | NameNode + ResourceManager + Spark Master |
| 4个工作节点 | 8核 / 32GB内存 | DataNode + NodeManager + Spark Executor运行 |
| 存储 | HDFS 3副本 | 数据冗余 |
这个配置并不高,但关键问题是:跑Spark作业时,executor核心数和内存如何分配?这决定了你能拿到多少并行度。总资源就这么多,分配给每个Executor的核数越少,Executor数量越多,并行度越高;但任务之间的网络Shuffle成本也可能上升。实操中,常见的经验值是单个Executor给2~4个core,内存按容器内预留一点给YARN overhead后尽量都给Spark使用。
2.2 真正影响性能的部署参数
安装完Spark之后,先别急着跑作业,有几个配置参数是必须过一遍的,因为它们直接决定资源利用率。以下参数在spark-defaults.conf里设置,或提交任务时用--conf传。
spark.cores.max 8 # 整个应用最多可用的CPU核心数 spark.executor.memory 8g # 每个Executor的堆内内存 spark.executor.memoryOverhead 2g # 每个Executor的堆外内存 spark.executor.cores 2 # 每个Executor占用的CPU核心数 spark.sql.shuffle.partitions 200 # Spark SQL Shuffle默认分区数 spark.default.parallelism 8 # RDD默认分区数 spark.serializer org.apache.spark.serializer.KryoSerializer很多初学的人会以为spark.executor.memory给得越大越好,实际并不完全是。Executor在YARN上申请的内存如果太大,JVM的GC会成为新瓶颈。我见到过不少例子,把Executor内存调到32G甚至64G,结果Full GC频繁,整个作业运行期间GC时间占到了30%以上。一般建议单个Executor堆内内存在8G到16G之间,对绝大多数批处理作业来说已经足够。
KryoSerializer值得单独说。Spark默认的Java序列化器方便但性能很差。对于数据量大、或者使用了自定义类的场景,用spark.serializer切换Kryo能够明显减少序列化时间和内存占用。有些网上教程会建议在spark-defaults.conf里全局开启Kryo,但如果你的代码里注册了自定义类,一定要用spark.kryo.registrator或spark.kryo.classesToRegister提前注册,否则Kryo在遇到未注册类时仍然会走落伍的路径,性能优势大打折扣。
3. Shuffle调优:MapReduce和Spark的共同命门
3.1 MapReduce侧怎么调Shuffle
很多人从MapReduce迁移到Spark后,容易忽略一个事实:Shuffle仍然是两个框架共同的性能命门。MapReduce的Shuffle包括Map输出端的Sort和Spill,以及Reduce端从各节点拉取中间结果的Copy。这里的Sort到底能不能关?在Hadoop早期的版本里,如果想在Map端做Combine(比如本地聚合),必须先Sort保证Key有序才能高效合并,所以Sort几乎是强制性的。后来Hadoop做了优化,有一种Secondary Sort机制,但那时写Java代码的复杂度已经让很多人放弃了。
在MapReduce里实际能调的Shuffle参数主要这几项:mapreduce.task.io.sort.mb决定Map输出缓冲区大小,改大了能减少Spill次数;mapreduce.reduce.shuffle.parallelcopies决定Reduce端并行拉取Map输出的线程数,设得大能加快Copy阶段,但网络开销也会增加;mapreduce.map.sort.spill.percent默认0.8,触发Spill的阈值,调大也能减少落盘次数。这些参数作用都是缓解同一个老大难问题:数据搬运成本。
反过来我们会发现,Spark之所以能在Shuffle上更高效,是因为Spark引入了HashShuffle和后来的SortShuffle,经历过多个版本的迭代。到了Spark 2.x之后,默认的Shuffle方案是SortShuffleManager。它把Map的输出按照Partition写入内存缓冲,再排序落盘,每个Map任务最终生成的数据文件数量不再是“Reduce分区数”那种爆炸式增长,而是尽量合并,让Reduce端拉取更容易。
3.2 Spark侧Shuffle优化和数据倾斜处理
落到实际调优场景,Spark的Shuffle优化主要关注三点:分区数、缓冲大小、以及数据倾斜。
spark.sql.shuffle.partitions是Spark SQL作业里出现频率最高的一个参数。理论上说,Shuffle分区数越多,每个Task处理的数据量越小,并行度越高。但分区数太多,也会产生大量小的Shuffle文件,反而让调度和网络传输开销变大。合理区间建议是:Shuffle数据量在几十GB以内时,分区数设置在200到500之间通常表现不错;如果单个Task处理时间差异巨大,再看是否要动态开启AQE。
数据倾斜是Spark作业最常见的性能杀手。表现是:集群上大部分Task快速跑完,一两个Task卡了很久,最后整个Job等那个Task。解决倾斜没有万能法,常用的手段如下:
先定位倾斜的Key到底是什么。可用df.groupBy("key").count().orderBy(desc("count")).show()方式快速看分布。如果是热点Key数量不多,可以考虑加随机前缀再聚合,分两步做:第一步给热点Key打散,第二步去掉前缀归并真实结果。这个办法对两阶段聚合场景效果比较明显。
如果倾斜是发生在Join阶段,比如一张小表和一张大表Join,可以把小表用BroadcastHashJoin广播,避免Shuffle。如果本身是大表Join大表,可以先过滤掉脏数据(比如空Key),再考虑对热点前缀加盐。
Spill和GC也是个容易被忽视的问题。Shuffle读阶段如果抛出类似Shuffle file cannot find的异常,大概率是Executor在Fetch时内存不够,把曾经落盘的Shuffle文件删掉导致。遇到这种问题,要增加spark.shuffle.memoryFraction对应的内存,或者在Spark 2.x之后直接调spark.memory.fraction,同时检查Executor是否频繁GC。
4. 内存管理:从Spark OOM说起
4.1 运行时内存模型拆解
Spark作业跑挂的重要原因,除了资源没申请够之外,很多都和内存模型理解不到位有关。Spark的Executor内存可以拆成三块:堆内执行内存(Execution Memory)、堆内存储内存(Storage Memory)、以及堆外内存(Off-heap Memory)。
在旧版本里,Execution Memory和Storage Memory是按比例静态划分的,两边的空间不能互相借用,导致经常出现“执行内存不够但存储内存闲置”的尴尬。从1.6开始引入统一内存管理(Unified Memory),让Execution和Storage可以互相抢占多余的闲置空间,这是一个非常大的改进。
但统一内存不代表没有OOM。通常说的Spark OOM,一部分是JVM堆内存真正不够,OutOfMemoryError: Java heap space;另一部分是堆外内存耗尽,也就是Container killed by YARN for exceeding memory limits,这个在日志里很常见。后者的起因往往是spark.executor.memoryOverhead或spark.memory.offHeap.size设置不合理。尤其在使用外部库(比如读取Redis、访问JDBC数据源)的时候,底层库分配的直接内存或线程栈有可能不计入Spark的堆内预算,一旦超过YARN容器阈值,就会被kill。
4.2 常见的OOM场景排查步骤
如果你发现Spark作业反复OOM,先不要急着一味加Executor内存,按下面的顺序排查更高效:
第一步,看是Driver端还是Executor端OOM。Driver端OOM通常是因为collect()把全量结果抓到内存,或广播变量太大。Executor端OOM则要去分析Task处理的数据量。
第二步,看日志确认堆内还是堆外。堆内OOM一般是数据分区过大或代码里有巨大的集合缓存;堆外OOM则多半是序列化缓冲或外部资源占用过多。
第三步,检查是否存在数据倾斜。如果一个处理流程中,某个Task拉取的数据量特别大,那么就是典型倾斜,把倾斜的Key打散或调大并行度,比简单加内存效率高得多。
这里我额外想强调一个亲身踩过的坑:别在RDD和DataFrame的算子内部维护大型集合对象。比如用mapPartitions做外部API调用时,把整个Partition的数据攒成ArrayList再批量提交,分区的数据量大时,这个List会直接在Executor堆内产生巨大压力。我当时就是在一个数据清洗作业里,用mapPartitions批量读取Redis的值,为了减少网络往返,把一个Partition内上万条Key的Value全部攒到List再写库,结果整个Executor被撑满,日志报的是堆内OOM。后来改成边读边写,按批大小限制在1000条以内,内存问题立刻消失。
5. Spark SQL性能优化与外部数据源集成实践
5.1 Spark SQL调优三板斧
实际生产环境里,大部分人写的不是RDD算子,而是Spark SQL。所以Spark SQL的性能调优更需要熟练掌握。我把它归纳成三板斧:代码本身的执行计划优化、资源配置与运行时参数优化、外部存储IO优化。
第一板斧是执行计划优化。Spark SQL通过Catalyst优化器生成物理执行计划,你要做的是学会用explain()去看执行计划,确认Join的方式是BroadcastHashJoin还是SortMergeJoin;确认有没有不必要的全表扫描,有没有可以做下推过滤却没做的情况。动态分区裁剪和文件中列裁剪这些能力在Spark 3.x里默认开启了一部分,但某些嵌套子查询还会退化成全量扫描,这时候可以通过改写SQL让过滤条件下推。
第二板斧是资源配置。AQE在Spark 3.x默认是开启的,Spark 3.4里相关参数比较完整,包括自动处理Join时数据倾斜、自动合并Shuffle分区、自动切换Join策略。AQE是好东西,但对于依赖稳定分区数的报表任务,也有可能出现“动态合并后某个大分区还是过重”的情况。出现这种问题时,建议固定spark.sql.adaptive.coalescePartitions.enabled为false,回到手动设置分区数的方式。
第三板斧是外部存储IO优化。这是Spark SQL实战里最容易忽略的地方,也是最容易出效果的地方。以Spark集成Redis为例,很多人直接写一个UDF,在Task内部遍历调用Redis客户端获取数据。这种做法如果键数量少还好,键一多,网络往返延迟就会被放得很大。较合理的方案是使用Redis的pipeline或批量获取命令;或者将需要关联的数据预读并广播到每个Executor,变成本地内存查找。当然,如果预读数据很大,就老老实实做Join。
Spark SQL读取外部数据库也有类似道理。拿达梦数据库(DM)来举例,国产化环境中需要对接到Spark时,应该怎么处理?重点之一是JDBC连接参数。每次读取数据,numPartitions决定并行度,lowerBound、upperBound决定拆分列范围。如果一张千万级的大表只用一个连接去拉,Spark会先给你一个Timeout。给JDBC Reader配置合理的分区数量、每次批量拉取的行数、连接超时时间,能明显缓解全量拉取时的压力。另一个容易踩的点是底层驱动不一定支持下游的数据类型,读取时尽量用dbtable的方式写子查询,提前把列映射到Spark认识的类型上。
5.2 一个可参考的实际调优案例
把前面这些点穿起来看一个简单的分析案例,场景是:读取存放在HDFS上的用户行为明细表,需要先用Spark SQL做数据清洗,然后关联一份从Redis读取的白名单数据,最后关联达梦数据库中的用户维度表做宽表输出。
整个过程可能会遇到三个层面的问题:清洗阶段Shuffle太多、关联Redis太慢、读取达梦库形成单点瓶颈。对应调整的思路:
清洗阶段先过滤后聚合。提前用filter做行裁剪,再用select做列裁剪,减少进入Shuffle的数据量;Join白名单时,如果白名单本身是百万量级以内,用广播变量方式将白名单做成Map并在Task内直接lookup,而不是对每条明细发起Redis命令。接入达梦维度表时,给JDBC源设置合理的分区列,尽量使用主键或数值型字段做partitionColumn,同时按Where条件切分多个子查询并发读取,避免一个连接拉全表。
调整之后,整个任务从原来的20多分钟缩短到5分钟以内。这个效果并不来自某一种高深的魔法,而是每一个IO环节都减掉了多余的搬运和等待。
6. 数据倾斜之外的一些“软性”调优点
6.1 从MapReduce编程实例里学到的资源利用思维
处理过一个传统MapReduce编程实例:统计网站上每个页面的PV和UV。这个实例有多个阶段,需要两个MR作业串联:第一个MR完成PV统计并输出按UV去重后的预聚合数据,第二个MR做最终汇总。如果严格按照两个独立的MapReduce作业来跑,每个作业都要把中间结果写一遍HDFS再拉一遍,整个链路的时间翻倍。后来我们在第一个Reducer里把预聚合做得足够细致,让第二个Job的Map阶段几乎只做读取,全局排序的意义已经不大,这样总耗时下降不少。
这个思路放到Spark里同样适用。不要因为Spark支持在DAG里做多Stage就随意把所有操作都串联起来,每个Stage之间的宽依赖仍然对应一次Shuffle。写代码之前先在纸上画出数据流,标注哪些地方会发生Shuffle,哪些宽依赖其实可以转化成窄依赖。比如,能用reduceByKey做局部聚合再用全局聚合的地方,就别一上来直接groupByKey。这个原则从MapReduce时代到Spark时代都没有变过。
6.2 Spark面试常见考点对应的调优理解
顺着话题多说几句,Spark面试题里经常出现“Spark为什么比MapReduce快”、“数据倾斜怎么解决”、“Spark OOM怎么处理”、“Spark SQL的执行流程是什么”。这些考点看起来分散,其实对应的都是一个工程师做调优时需要掌握的基本盘:快在计算模型和调度模型,数据倾斜是Shuffle不均,OOM是内存模型没吃透,Spark SQL是在Catalyst优化器之上写SQL。如果能把这几个点想透,面试和实际调优基本上都没有太大障碍。反过来说,如果只是一味背参数,而不知道参数背后的模型,换一个集群环境就会手足无措。
有个新趋势是DGX Spark这类硬件加速集成方案开始出现,它把Spark的批处理跑在GPU 加速硬件上,这类方案要求对执行模型的IO瓶颈和调度开销同样有比较深的理解,否则也是白搭。遇到这种环境,建议从小的数据分析案例入手,先在无状态作业里跑通,再逐步加复杂依赖,否则很难定位新增的瓶颈到底在算子还是调度层。
6.3 集群日常维护的调优习惯
最后补一些日常操作细节。Spark集群搭好,作业也能跑了,并不意味着可以躺平。我自己习惯在每类任务跑完后去Spark UI上面看几个数值:Executor的GC时间、Shuffle读写量、单个Task耗时分布、以及是否频繁发生Speculative Task重试。这些数据比任何理论都更能说明问题。如果GC时间偏高,优先处理代码中多余的缓存对象;如果Shuffle数据量远远大于输入数据量,说明在Shuffle上游可以做更彻底的预聚合或列裁剪;如果Task耗时中位数很低但某些尾巴特别高,大概率又是倾斜问题。
还有每周清理一次Spark日志和Event Log的习惯,防止磁盘占用过高导致节点异常。看起来是运维活,但生产环境的速度和稳定性往往来自这些不起眼的地方。
我在实际使用中也发现,很多人把性能调优理解成“加资源、加并行度”,没解决本质的Shuffle和IO问题时,资源堆上去只是暂时掩盖问题,数据量一涨又被打回原形。真正有效的调优顺序永远是:先定位瓶颈在哪一层,再对该层做最便宜的优化。MapReduce到Spark的迁移只是提供了更多优化的空间,不是终点。遇到新框架、新硬件平台时,这个思路依然管用。