作者信息

作者: 余辉       微信公众号:辉哥大数据
购买地址: 京东淘宝当当网
读者须知:本书配套示例源码、PPT课件、教学视频与作者答疑服务,购买之后可加粉丝群

本书封面

在这里插入图片描述


第13章 Spark性能调优

       本章将全面剖析Spark性能调优的方法,从常规调优到开发原则优化,再到具体的调优方法和数据倾斜问题的应对策略,层层递进,帮助读者深入理解并掌握Spark调优的精髓。通过本章的学习,读者将能够显著提高调优Spark的性能。

       本章主要知识点:

  • Spark常规性能调优
  • Spark开发原则优化
  • Spark调优方法
  • Spark数据倾斜调优

13.1 Spark常规性能调优

       本节将深入探讨Spark的常规性能调优策略,涵盖最优资源配置、RDD优化、并行度调节、广播大变量、Kryo序列化、调节本地化等待时长、ShuGle调优及JVM调优等多种方法,旨在帮助用户全面优化Spark作业性能。

13.1.1 常规性能调优一:最优资源配置

       Spark性能调优的第一步就是为任务分配更多的资源。在一定范围内,增加资源的分配与性能的提升是成正比的,实现最优的资源配置后,在此基础上再考虑进行后面讲解的性能调优策略。资源的分配在使用脚本提交Spark任务时进行指定,标准的Spark任务提交脚本代码如下:

/usr/opt/modules/spark/bin/spark-submit \
--class com.atguigu.spark.Analysis \
--num-executors 80 \
--driver-memory 6g \
--executor-memory 6g \
--executor-cores 3 \
/usr/opt/modules/spark/jar/spark.jar 

可以进行分配的资源如表13-1所示。
在这里插入图片描述
       调节原则:尽量将任务分配的资源调节到可以使用的资源的最大限度。对于具体资源的分配,我们分别讨论Spark的两种Cluster运行模式:

  • 第一种是Spark Standalone模式。在提交任务前,一定要知道或者可以从运维部门获取到可以使用的资源情况。在编写Submit脚本时,应根据可用的资源情况进行资源的分配。例如集群有15台机器,每台机器有8GB内存和两个CPU Core,那么可以指定15个Executor,每个Executor分配8GB内存和两个CPU Core。
  • 第二种是Spark YARN模式。由于YARN使用资源队列进行资源的分配和调度,在编写Submit脚本时,应根据Spark作业要提交到的资源队列进行资源分配。例如,资源队列有400GB内存和100个CPU Core,那么指定50个Executor,每个Executor分配8GB内存和2个CPU Core。

13.1.2 常规性能调优二:RDD优化

针对RDD的优化有3种措施,即RDD复用、RDD持久化和RDD尽可能早的Filter操作。
1)RDD复用
       获取到初始RDD后,需要检查相同的算子和计算逻辑,避免相同计算逻辑的重复执行。
2)RDD持久化
       在Spark中,当多次对同一个RDD执行算子操作时,每一次都会对这个RDD以之前的父RDD重新计算一次,这种情况是必须要避免的,因此对同一个RDD的重复计算是对资源的极大浪费。因此,必须对多次使用的RDD进行持久化,通过持久化将公共RDD的数据缓存到内存/磁盘中,之后对于公共RDD的计算都会从内存/磁盘中直接获取RDD数据。

对于RDD的持久化,有两点需要说明:
(1)RDD的持久化是可以进行序列化的,当内存无法完整地存放RDD的数据时,可以考虑使用序列化的方式减小数据体积,将数据完整存储在内存中。
(2)如果对于数据的可靠性要求很高,并且内存充足,可以使用副本机制对RDD数据进行持久化。当启用复本机制时,对于持久化的每个数据单元都会在其他节点上存储一个副本,由此实现数据的容错。一旦一个副本数据丢失,不需要重新计算,还可以使用另一个副本。
3)RDD尽可能早的Filter操作
获取到初始RDD后,应该考虑尽早过滤掉不需要的数据,进而减少对内存的占用,从而提升Spark作业的运行效率。

13.1.3 常规性能调优三:并行度调节

       Spark作业中的并行度是指各个Stage中Task的数量。如果并行度设置不合理,导致并行度过低,就会导致资源的极大浪费。例如,如果有20个Executor,每个Executor分配3个CPU Core,而Spark作业有40个Task,这样每个Executor分配到的Task个数是两个。这就使得每个Executor有一个CPU Core空闲,导致资源的浪费。
       理想的并行度设置应该是让并行度与资源相匹配。简单来说,就是在资源允许的前提下,并行度要设置得尽可能大,以充分利用集群资源。合理设置并行度可以提升整个Spark作业的性能和运行速度。
       Spark官方推荐Task数量应该设置为Spark作业总CPU Core数量的2~3倍。之所以不推荐Task数量与CPU Core总数相等,这是因为Task的执行时间不同,有的Task执行速度快,而有的Task执行速度慢。如果Task数量与CPU Core总数相等,那么执行快的Task执行完成后,会出现CPU Core空闲的情况。如果Task数量设置为CPU Core总数的2~3倍,那么一个Task执行完毕后,CPU Core会立刻执行下一个Task,从而减少资源浪费,同时提升Spark作业运行的效率。

13.1.4 常规性能调优四:广播大变量

       默认情况下,Task的算子中如果使用了外部变量,每个Task都会获取一份变量的复本,这就造成了内存的极大消耗。一方面,如果后续对RDD进行持久化,可能就无法将RDD数据存入内存,只能写入磁盘,磁盘I/O将会严重消耗性能;另一方面,Task在创建对象时,也许会发现堆内存无法存放新创建的对象,这就会导致频繁的GC,GC会导致工作线程停止,进而导致Spark暂停工作一段时间,严重影响Spark的性能。
       假设当前任务配置了20个Executor,指定500个Task,有一个20MB的变量被所有Task共用,此时会在500个Task中产生500个副本,耗费集群10GB的内存。如果使用了广播变量,那么每个Executor保存一个副本,一共消耗400MB内存,内存消耗减少了1/5。
       广播变量在每个Executor中只保存一个副本,此Executor的所有Task共用此广播变量,这使得变量产生的副本数量大大减少。
       在初始阶段,广播变量只在Driver中有一份副本。Task在运行时,想要使用广播变量中的数据,此时首先会在自己本地的Executor对应的BlockManager中尝试获取变量。如果本地没有这个变量,BlockManager就会从Driver或者其他节点的BlockManager上远程拉取变量的复本,并由本地的BlockManager进行管理;之后此Executor的所有Task都会直接从本地的BlockManager中获取变量。

13.1.5 常规性能调优五:Kryo序列化

       默认情况下,Spark使用Java的序列化机制。Java的序列化机制使用方便,不需要额外的配置,在算子中使用的变量实现Serializable接口即可。但是,Java序列化机制的效率不高,序列化速度慢,并且序列化后的数据占用的空间依然较大。
       Kryo序列化机制比Java序列化机制性能提高10倍左右。Spark之所以没有默认使用Kryo作为序列化类库,是因为它不支持所有对象的序列化;同时,Kryo需要用户在使用前注册需要序列化的类型,不够方便。但从Spark 2.0.0版本开始,简单类型、简单类型数组、字符串类型的ShuGling RDD已经默认使用Kryo序列化方式了。

13.1.6 常规性能调优六:调节本地化等待时长

       Spark作业运行过程中,Driver会对每一个Stage的Task进行分配。根据Spark的Task分配算法,Spark希望Task能够运行在它要计算的数据所在的节点(数据本地化思想),这样就可以避免数据的网络传输。通常来说,Task可能不会被分配到它处理的数据所在的节点,因为这些节点可用的资源可能已经用尽。此时,Spark会等待一段时间,默认为3s,如果等待指定时间后仍然无法在指定节点运行,那么会自动降级,尝试将Task分配到比较差的、本地化级别对应的节点上,比如将Task分配到离它要计算的数据比较近的一个节点,然后进行计算,如果当前级别仍然不行,那就继续降级。
       当Task要处理的数据不在Task所在节点上时,会发生数据的传输。Task会通过所在节点的BlockManager获取数据,BlockManager发现数据不在本地时,会通过网络传输组件从数据所在节点的BlockManager处获取数据。
       网络传输数据的情况是我们不愿意看到的,大量的网络传输会严重影响性能。因此,我们希望通过调节本地化等待时长,如果在等待时长这段时间内,目标节点处理完成了一部分Task,那么当前的Task将有机会得到执行,这样就能够改善Spark作业的整体性能。
       在Spark项目开发阶段,可以使用Client模式对程序进行测试。此时,可以在本地看到比较全的日志信息。日志信息中有明确的Task数据本地化的级别,如果大部分都是PROCESS_LOCAL,那么无须进行调节。但是,如果发现很多级别都是NODE_LOCAL、ANY,那么就需要对本地化的等待时长进行调节,通过延长本地化等待时长,看看Task的本地化级别有没有提升,并观察Spark作业的运行时间有没有缩短。
       注意,过犹不及,不要将本地化等待时长延长得过长,导致因为大量的等待时长,使得Spark作业的运行时间反而增加了。
       Spark本地化等待时长的设置代码如下:

val conf = new SparkConf().set("spark.locality.wait", "6") 

13.1.7 常规性能调优七:ShuGle调优

       1)ShuGle调优一:调节Map端缓冲区大小
       在Spark任务运行过程中,如果Shuffle的Map端处理的数据量比较大,但是Map端缓冲的大小是固定的,可能会出现Map端缓冲数据频繁溢写到磁盘文件中的情况,使得性能非常低下。通过调节Map端缓冲的大小,可以避免频繁的磁盘I/O操作,进而提升Spark任务运行的整体性能。
       Map端缓冲的默认配置是32KB,如果每个Task处理640KB的数据,那么会发生640/32= 20次溢写;如果每个Task处理64000KB的数据,就会发生64000/32=2000次溢写,这对于性能的影响是非常严重的。
       Map端缓冲的配置代码如下:

val conf = new SparkConf().set("spark.shuffle.file.buffer", "64")

       2)ShuGle调优二:调节Reduce端拉取数据缓冲区大小
       在Spark Shuffle过程中,Shuffle Reduce Task的缓冲区大小决定了Reduce Task每次能够缓冲的数据量,也就是每次能够拉取的数据量。如果内存资源较为充足,适当增加拉取数据缓冲区的大小,可以减少拉取数据的次数,从而减少网络传输的次数,进而提升性能。
       Reduce端数据拉取缓冲区的大小可以通过spark.reducer.maxSizeInFlight参数进行设置,默认为48MB。该参数的设置代码如下:

val conf = new SparkConf().set("spark.reducer.maxSizeInFlight", "96")

       3)ShuGle调优三:调节Reduce端拉取数据重试次数
       在Spark Shuffle过程中,Reduce Task拉取属于自己的数据时,如果因为网络异常等原因导致失败,会自动进行重试。对于那些包含特别耗时的Shuffle操作的作业,建议增加重试最大次数(比如60次),以避免由于JVM的Full GC或者网络不稳定等因素导致的数据拉取失败。在实践中发现,对于针对超大数据量(数十亿到上百亿)的Shuffle过程,调节该参数可以大幅提升稳定性。
       Reduce端拉取数据重试次数可以通过spark.shuule.io.maxRetries参数进行设置,该参数就代表了可以重试的最大次数。如果在指定次数之内拉取还是没有成功,可能会导致作业执行失败。该参数默认为3,其设置代码如下:
val conf = new SparkConf().set(“spark.shuffle.io.maxRetries”, “6”)
       4)ShuGle调优四:调节Reduce端拉取数据等待间隔
在Spark Shuffle过程中,Reduce Task拉取属于自己的数据时,如果因为网络异常等原因导致失败会自动进行重试,在一次失败后,会等待一定的时间间隔再进行重试,这样可以通过加大间隔时长(比如60s)来增加Shuffle操作的稳定性。
       Reduce端拉取数据等待间隔可以通过spark.shuffle.io.retryWait参数进行设置,默认值为5s。该参数的设置代码如下:

val conf = new SparkConf().set("spark.shuffle.io.retryWait", "60s")

       5)ShuGle调优五:调节SortShuGle排序操作阈值
       对于SortShuffleManager,如果Shuffle Reduce Task的数量小于某一阈值,则Shuffle write过程中不会进行排序操作,而是直接按照未经优化的HashShuffleManager方式来写数据,但是最后会将每个Task产生的所有临时磁盘文件都合并成一个文件,并会创建单独的索引文件。
       当你使用SortShuffleManager时,如果确实不需要排序操作,那么建议将这个参数调大一些,大于Shuffle Read Task的数量。这样,map-side就不会进行排序,从而减少排序的性能开销。但在这种方式下,依然会产生大量的磁盘文件,因此Shuffle write性能有待提高。
       SortShuffleManager排序操作阈值可以通过spark.shuffle.sort.bypassMergeThreshold这一参数进行设置,默认值为200。该参数的设置代码如下:

val conf = new SparkConf()
.set("spark.shuffle.sort.bypassMergeThreshold", "400")

13.1.8 常规性能调优八:JVM调优

       对于JVM调优,首先应该明确,full GC/minor GC都会导致JVM的工作线程停止工作,即stop the world。本小节将分析JVM调优的几种方法。
       1. JVM调优一:降低Cache操作的内存占比
       1)静态内存管理机制
       根据Spark静态内存管理机制,堆内存被划分为两部分:Storage和Execution。Storage主要用于缓存RDD数据和广播数据,Execution主要用于缓存在Shuule过程中产生的中间数据。Storage占系统内存的60%,Execution占系统内存的20%,并且两者完全独立。
       在一般情况下,Storage的内存都提供给Cache操作,但是如果在某些情况下Cache操作内存不是很紧张,而Task的算子中创建的对象很多,Execution内存又相对较小,这会导致频繁的minor GC,甚至full GC,进而导致Spark频繁地停止工作,对性能的影响会很大。
       在Spark UI中可以查看每个Stage的运行情况,包括每个Task的运行时间、GC时间等。如果发现GC太频繁,时间太长,就可以考虑调节Storage的内存占比,让Task执行算子函数时有更多的内存可以使用。
       Storage内存区域可以通过spark.storage.memoryFraction参数进行指定,默认为0.6,即60%,可以逐级向下递减,其配置代码如下:

val conf = new SparkConf().set("spark.storage.memoryFraction", "0.4") 

       2)统一内存管理机制
       根据Spark统一内存管理机制,堆内存被划分为两部分:Storage和Execution。Storage主要用于缓存数据,Execution主要用于缓存在Shuffle过程中产生的中间数据,两者所组成的内存部分称为统一内存,Storage和Execution各占统一内存的50%。由于动态占用机制的实现,Shuffle过程需要的内存过大时,会自动占用Storage的内存区域,因此无须手动进行调节。

       2. JVM调优二:调节Executor堆外内存
       Executor的堆外内存主要用于程序的共享库、Perm Space、线程Stack和一些Memory Mapping等,或者类C方式分配对象。
       有时,如果你的Spark作业处理的数据量非常大,达到几亿的数据量,此时运行Spark作业会时不时地报错,例如Shuffle output file cannot find、Executor lost、Task lost、Out of memory等。这可能是Executor的堆外内存不太够用,导致Executor在运行过程中内存溢出。
       Stage的Task在运行时,可能要从一些Executor中拉取Shuffle Map Output文件。但是Executor可能已经由于内存溢出挂掉,其关联的BlockManager也会丢失,这就可能会报出Shuffle output file cannot find、Executor lost、Task lost、Out of memory等错误。此时可以考虑调节一下Executor的堆外内存,也就可以避免报错。与此同时,堆外内存调节得比较大时,对于性能来讲,也会带来一定的提升。
       默认情况下,Executor堆外内存上限大概多于300MB。在实际生产环境下,对海量数据进行处理时,这里往往会出现问题,导致Spark作业反复崩溃,无法运行。此时,可以将该参数调节到至少1GB,甚至2GB、4GB。
       Executor堆外内存的配置需要在spark-submit脚本中进行设置,配置语句如下:

--conf spark.yarn.executor.memoryOverhead=2048 

       以上参数配置完成后,会避免某些JVM OOM的异常问题,同时提升整体Spark作业的性能。
3. JVM调优三:调节连接等待时长
       在Spark作业运行过程中,Executor优先从自己本地关联的BlockManager中获取某份数据。如果本地BlockManager没有数据,会通过TransferService远程连接其他节点上Executor的BlockManager来获取数据。
       如果Task在运行过程中创建大量对象或者创建的对象较大,会占用大量的内存,这会导致频繁的垃圾回收。但是垃圾回收会导致工作现场全部停止,也就是说,垃圾回收一旦执行,Spark的Executor进程就会停止工作,无法提供响应。此时,由于没有响应,无法建立网络连接,会导致网络连接超时。
       在生产环境下,有时会遇到file not found、file lost这类错误。在这种情况下,很有可能是Executor的BlockManager在拉取数据时无法建立连接,然后在超过默认的连接等待时长60s后,宣告数据拉取失败。如果反复尝试都拉取不到数据,可能会导致Spark作业的崩溃。这种情况也可能会导致DAGScheduler反复提交几次Stage,TaskScheduler反复提交几次Task,大大延长了Spark作业的运行时间。
       此时,可以考虑调节连接的超时时长,连接等待时长需要在spark-submit脚本中进行设置,设置方式如下:

--conf spark.core.connection.ack.wait.timeout=300 2

调节连接等待时长后,通常可以避免部分的某某文件拉取失败、某某文件lost等报错。

13.2 Spark开发原则优化

       本节将阐述Spark开发中的关键优化原则,包括避免创建重复的RDD与DataFrame、复用RDD以提高效率、减少重复性SQL查询、注意数据类型选择以及编写高质量的SQL语句等,旨在帮助开发者提升Spark应用的性能与代码质量。

13.2.1 开发原则一:避免创建重复的RDD

       通常来说,我们在开发一个Spark作业时,首先是基于某个数据源(比如Hive表或HDFS文件)创建一个初始的RDD,接着对这个RDD执行某个算子操作,然后得到下一个RDD,以此类推,循环往复,直到计算出最终我们需要的结果。在这个过程中,多个RDD会通过不同的算子操作(比如Map、Reduce等)串起来,这个“RDD串”就是RDD lineage,也就是“RDD的血缘关系链”。
       在开发过程中要注意:对于同一份数据,只应该创建一个RDD,不能创建多个RDD来代表同一份数据。一些Spark初学者在刚开始开发Spark作业时,或者是有经验的工程师在开发RDD lineage极其冗长的Spark作业时,可能会忘记自己之前已经为某份数据创建过一个RDD,从而导致为同一份数据创建了多个RDD。这就意味着,我们的Spark作业会进行多次重复计算,来创建多个代表相同数据的RDD,进而增加了作业的性能开销。
       下面举例说明。我们需要对名为hello.txt的HDFS文件进行一次Map操作,再进行一次Reduce操作。也就是说,需要对一份数据执行两次算子操作。

// 错误的做法:对同一份数据执行多次算子操作时,创建多个RDD
val rdd1 = sc.textFile("hdfs:// 192.168.0.0:8020/hello.txt")
rdd1.map(...)
val rdd2 = sc.textFile("hdfs:// 192.168.0.0:8020/hello.txt")
rdd2.reduce(...)

       这里执行了两次textFile方法,即针对同一个HDFS文件创建了两个RDD,然后分别对每个RDD执行一个算子操作。在这种情况下,Spark需要从HDFS上加载两次hello.txt文件的内容,并创建两个单独的RDD。第二次加载HDFS文件以及创建RDD的性能开销是明显浪费的。

// 正确的做法:对同一份数据执行多次算子操作时,只使用一个RDD
val rdd1 = sc.textFile("hdfs:// 192.168.0.0:8020/hello.txt")
rdd1.map(...)
rdd1.reduce(...)

       这种写法很明显比上一种写法好多了,因为我们对于同一份数据只创建了一个RDD,然后对这个RDD执行了多次算子操作。但要注意,到这里优化还没有结束,由于rdd1被执行了两次算子操作,第二次执行Reduce操作时,还会再次从源头处重新计算一次rdd1的数据,因此还是会有重复计算的性能开销。
       要彻底解决这个问题,必须结合后续讲解的“开发原则三:尽可能复用同一个RDD”,才能保证一个RDD在多次使用时只被计算一次。

13.2.2 开发原则二:避免创建重复的DataFrame

       重复创建相同的DataFrame是初学者比较容易犯的一个错误。
对于一个会被多次使用的数据集,我们应该只创建一个DataFrame实例来表示它。很多读者写代码时有复制粘贴的习惯,不是说这个习惯不好,而是这个习惯很容易导致相同的DataFrame在无意中被创建多次。如果对同一个数据集创建了多个相同的DataFrame实例,就会浪费内存资源,甚至还会导致重复计算的问题。因此,读者在写代码时一定要留心,避免重复创建相同的DataFrame实例。
       下面是一个低效的代码例子:

val spark = SparkSession.builder().getOrCreate()
...
for( a <- 1 to 10){
  val df=spark.read.json(“employee.json”)
  ...
}
...

       上述代码在for循环中反复创建了相同的DataFrame,这对资源造成了浪费。注意,要避免编写这种低效率的代码。
       修改方法是把df放到for循环外面。修正后的代码如下:

val spark = SparkSession.builder().getOrCreate()
...
val df = spark.read.json(“employee.json”)
for( a <- 1 to 10){
  ...
}
...

13.2.3 开发原则三:尽可能复用同一个RDD

       除了要避免在开发过程中对同一份数据创建多个RDD外,在对不同的数据执行算子操作时,还要尽可能地复用同一个RDD。例如,有一个RDD的数据格式是key-value类型的,另一个是单value类型的,这两个RDD的value数据完全一样,那么此时可以只使用key-value类型的RDD,因为其中已经包含另一个RDD的数据。对于类似这种多个RDD的数据有重叠或者包含的情况,我们应该尽量复用一个RDD,这样可以尽可能减少RDD的数量,从而尽可能减少算子执行的次数。
       下面来看一个例子。以下代码中有一个<Long, String>格式的RDD,即rdd1,由于业务需要,对rdd1执行了一个Map操作,创建了一个rdd2,而rdd2中的数据仅仅是rdd1中的value值,也就是说,rdd2是rdd1的子集。

JavaPairRDD<Long, String> rdd1 = ...
JavaRDD<String> rdd2 = rdd1.map(...)

// 分别对rdd1和rdd2执行不同的算子操作
rdd1.reduceByKey(...)
rdd2.map(...)

       在这个例子中,rdd1和rdd2只是数据格式不同,rdd2的数据完全是rdd1的子集,却创建了两个RDD,并对两个RDD都执行了一次算子操作。此时会因为对rdd1执行Map算子来创建rdd2,而多执行一次算子操作,进而增加性能开销。其实在这种情况下,完全可以复用同一个RDD。我们可以使用rdd1既进行reduceByKey操作,又进行Map操作。在进行第二个Map操作时,只使用每个数据的tuple._2,也就是rdd1中的value值即可。

JavaPairRDD<Long, String> rdd1 = ...
rdd1.reduceByKey(...)
rdd1.map(tuple._2...)

       第二种方式相较于第一种方式而言,很明显减少了一次rdd2的计算开销。但是到这里优化还没有结束,对rdd1还是执行了两次算子操作,rdd1实际上还是会被计算两次。因此,还需要配合“对多次使用的RDD进行持久化”,才能保证一个RDD在多次使用时只被计算一次。

13.2.4 开发原则四:避免重复性的SQL查询,对DataFrame复用

这个用语言表述比较麻烦,我们直接看代码吧。
假设有一张students表,用以下示例代码对这张表的操作是低效的:

val spark = SparkSession.builder().getOrCreate()
import spark.implicits._
val studentNameAndAge = spark.sql(“SELECT name,age FROM students WHERE class=1)
...
// 经过多行代码之后
...
val studentName = spark.sql(“SELECT name FROM students WHERE class=1 AND age > 20)

上面这段代码对students表查询了两次,如果这张表特别大,查询的效率就会很低。
这里讲一下,在Spark的Scala API中,DataFrame的定义是这样的:

type DataFrame = Dataset[Row]

因此,DataFrame可以使用Dataset的一些方法。
下面是修正之后的代码:

val spark = SparkSession.builder().getOrCreate()
import spark.implicits._
val studentNameAndAge = spark.sql(“SELECT name,age FROM students WHERE class=1)
...
// 经过多行代码之后
...
val studentName = studentNameAndAge.filter($"age" > 20)

       这样代码的效率就会提高。因为开发程序时通常会写很多行代码,许多人写着写着就忘记了前面的SQL和接下来的SQL是否执行了相同效果的查询(代码一长难免会不记得),导致了对表的不必要的重复查询。读者写代码时需要注意这个问题,避免写出低效的代码,尽量复用前面定义的DataFrame。

13.2.5 开发原则五:注意数据类型的使用

       这条原则有以下两点需要注意的地方。
       1)在生成DataFrame或Dataset时如何定义数据的类型
这个问题需要在我们定义变量时考虑清楚。Scala提供了丰富的数据类型,我们要根据场景选择合适的数据类型。能用Byte类型,就不要为了方便定义成Int类型。一个Byte类型是8位,而Int类型是32位,一旦将数据进行缓存,内存的消耗将会翻倍。在使用Spark SQL时,定义合适的数据类型可以节省比较可观的内存资源。
       2)在代码中能用基本类型,就尽量使用基本类型
       由于每个不同的Java对象都有一个“对象头”,这个“头”大概是16字节,里面包含一些信息,例如一个指向类的指针。
像String类型比Char类型的数组开销大40字节,因为它不仅存储了数据本身,还包含其他数据(比如String的长度)。同时,因为它是用UTF-16编码的,所以每个字符占2字节。因此,一个10字符长的String会占60字节的大小。
要避免使用类与对象中包含对象以及指针的这种嵌套结构,还可以考虑使用数值型的ID或者枚举类型替代用字符串表示的键。此外,常用的集合类,比如HashMap、LinkedList,使用链表的数据结构。其中每个元素(比如Map.Entry)都有一个“包装”对象,这种对象不仅具有前面所讲的“对象头”,还包含指向下一个对象的指针。
因此,对于自定义的对象、String、集合等,在不影响代码的可读性、可维护性的情况下,能不用就尽量不用,因为它们占用比较多的内存。

13.2.6 开发原则六:写出高质量的SQL

       在使用SQL查询时,一条高质量的SQL语句将节省大量的查询时间,以及节省宝贵的计算资源和内存资源。关于如何写出高质量的SQL语句,由于篇幅过长,也偏离了本书一开始定下的目标,这里就不展开描述了。
       如果读者之前接触过SQL的优化,想必听说过SQL的执行计划。获取执行计划是SQL优化很关键的一部分,接下来介绍一下如何获取SQL的执行计划。
这里建议读者在Spark Shell中执行以下语句,亲自编写语句能有效地理解语句的意思及其作用。
需要准备的数据:在HDFS中对应账户的文件夹下放置一个JSON文件(笔者将其放置在HDFS中的/spark_book_data/employee.json下),文件内容如下:

{"id" : "1201", "name" : "satish", "age" : "25"}
{"id" : "1202", "name" : "krishna", "age" : "28"}
{"id" : "1203", "name" : "amith", "age" : "39"}
{"id" : "1204", "name" : "javed", "age" : "23"}
{"id" : "1205", "name" : "prudvi", "age" : "23"}

准备的数据文件只有一个,下面我们进入Spark Shell。
(1)读取employee.json文件:

scala> val employee = spark.read.json("/spark_book_data/employee.json")
employee: org.apache.spark.sql.DataFrame = [age: string, id: string ... 1 more field]

(2)创建临时视图:

scala> employee.createOrReplaceTempView("employee")

(3)查看视图的Schema:

scala> employee.printSchema
root
 |-- age: string (nullable = true)
 |-- id: string (nullable = true)
 |-- name: string (nullable = true)

(4)通过toDebugString查看分区信息:

scala> employee.rdd.toDebugString
res2: String =
(1) MapPartitionsRDD[8] at rdd at <console>:24 []
 |  SQLExecutionRDD[7] at rdd at <console>:24 []
 |  MapPartitionsRDD[6] at rdd at <console>:24 []
 |  MapPartitionsRDD[5] at rdd at <console>:24 []
 |  FileScanRDD[4] at rdd at <console>:24 []

(5)获取SQL的执行计划。
在http:// :4040页面的SQL栏中,可以看到执行的详细情况,如图13-1所示。图13-2将显示的各种计划依次罗列了出来。
在这里插入图片描述
       最后总结一下,前面几个原则都是编写程序时需要注意的小细节,虽然看起来很简单,却能有效地提高代码的执行效率。不要嫌啰唆,因为细节决定成败。如果要处理的数据量很大,稍有不慎就会对时间以及计算资源造成极大的浪费。想要写出高质量的代码,仔细斟酌代码是非常有必要的。

13.3 Spark调优方法

       在大数据计算领域,Spark已经成为越来越流行、越来越受欢迎的计算平台之一。Spark的功能涵盖大数据领域的离线批处理、SQL类处理、流式/实时计算、机器学习、图计算等各种不同类型的计算操作,应用范围与前景非常广泛。然而,通过Spark开发出高性能的大数据计算作业并不是那么简单的。如果没有对Spark作业进行合理地调优,Spark作业的执行速度可能会很慢,这样就完全体现不出Spark作为一种快速大数据计算引擎的优势。因此,想要用好Spark,就必须对它进行合理的性能优化。Spark的性能调优实际上由很多部分组成,不会仅仅调节几个参数就可以立竿见影地提升作业性能。我们需要根据不同的业务场景以及数据情况对Spark作业进行综合性分析,然后进行多个方面的调节和优化,才能获得最佳性能。所有Spark作业都需要注意和遵循的一些基本原则,形成了较为常用的调优方法,这是高性能Spark作业的基础。本节将介绍几种Spark中的调优方法。

13.3.1 优化数据结构

在Java中,有3种数据类型比较耗费内存:
(1)对象,每个Java对象都有对象头、引用等额外的信息,因此比较占用内存空间。
(2)字符串,每个字符串内部都有一个字符数组以及长度等额外信息。
(3)集合类型,比如HashMap、LinkedList等,集合类型内部通常会使用一些内部类来封装集合元素,比如Map.Entry。
因此,Spark官方建议,在Spark编码实现中,特别是算子函数中的代码,尽量不要使用上述3种数据结构,尽量使用字符串替代对象,使用原始类型(比如Int、Long)替代字符串,使用数组替代集合类型,以尽可能地减少内存占用,从而降低GC频率,提升性能。
但是笔者在编码实践中发现,要做到该原则其实并不容易,因为我们同时要考虑代码的可维护性。如果一段代码中完全没有任何对象抽象,全部是字符串拼接的方式,那么对于后续的代码维护和修改,无疑是一场巨大的灾难。同理,如果所有操作都基于数组实现,而不使用HashMap、LinkedList等集合类型,那么对于编码难度以及代码的可维护性,也是一个极大的挑战。因此,笔者建议,在可能以及合适的情况下,使用占用内存较少的数据结构,但前提是要保证代码的可维护性。

13.3.2 使用缓存(Cache)

我们知道数据在内存中的计算非常快。因此,可以把需要进行多次操作的表缓存到内存中,避免对磁盘进行多次I/O操作。
缓存有两种方式,代码如下:

import org.apache.spark.storage._
val spark = SparkSession.builder().getOrCreate()
val df = spark.read.json(“employee”)
df.createOrReplaceTempView("employee")

// 方式一:缓存到内存中
spark.catalog.cacheTable(“employee”)  // 这样就缓存到内存中了
// 如果不需要缓存了,就清除它,清除方式如下
spark.catalog.uncacheTable("employee")
// 如果需要清除所有缓存,就使用clearCache()
spark.catalog.clearCache()
// 如果需要查看是否已经缓存,就使用isCached()
if( spark.catalog.isCached(“employee”) ){
	print("yes")
}else{
	print("no")
}
// 方式二:对DataFrame持久化

df.cache() 		// 这个默认的持久化级别是MEMORY_AND_DISK
df.persist() 	// 这个和上面一行代码的效果是一样的
df.persist(StorageLevel.MEMORY_ONLY)// 使用MEMORY_ONLY时的持久化级别
df.unpersist() 	// 释放

       除了手动进行缓存之外,Spark在执行Shuffle操作时也会自动对一些中间数据进行缓存,比如reduceByKey。
       上面第2种缓存方式是对DataFrame进行持久化。相比于RDD的默认持久化级别MEMORY_ONLY,DataFrame的默认持久化级别是MEMORY_AND_DISK。持久化级别及其说明如表13-2所示。
在这里插入图片描述

13.3.3 对配置属性进行调优

       在Spark 3.5.3版本中,Spark SQL的一些常见的优化配置属性如下。

  • spark.sql.files.maxPartitionBytes:控制每个分区的最大字节数,用于手动分区读取。
  • spark.sql.files.openCostInBytes:估算文件开销,用于手动分区读取。
  • spark.sql.broadcastTimeout:广播任务的超时时间。
  • spark.sql.shuffle.partitions:控制Shuffle过程中的分区数。
  • spark.sql.autoBroadcastJoinThreshold:控制自动广播合并的阈值大小。
  • spark.sql.codegen:是否启用代码生成优化。
  • spark.sql.costModel:选择不同的成本模型。
  • spark.sql.join.preferSortMergeJoin:是否优先考虑排序合并。
  • spark.sql.orc.filterPushdown:是否开启ORC格式下的过滤下推优化。
  • spark.sql.orc.vectorize:是否开启ORC格式下的矢量化读取优化。

       这些属性可以在Spark的配置文件spark-defaults.conf中设置,或者在创建SparkSession时通过.config()方法进行设置。例如:

val spark = SparkSession.builder()
  .appName("Spark SQL Optimization Example")
  .config("spark.sql.files.maxPartitionBytes", "67108864")
  .config("spark.sql.autoBroadcastJoinThreshold", "10485760")
  .getOrCreate()

注意,具体配置项可能随着Spark版本的更新而变化。
       需要说明的是,Spark每个版本的调优属性不一样,有新增的,也有删除的。要想获得最准确的信息,应当到官网找相应版本的文档进行查看。因此,这里只能简单地讲一下如何使用这些配置属性,具体的建议请读者查看官方文档Programming Guides中SQL里面的Performance Tuning小节。
       下面简单介绍多个版本中Spark SQL的一些参数的作用和区别。

  • spark.sql.inMemoryColumnStorage.compressed:默认值为true,它的作用是自动对内存中的列式存储进行压缩。
  • spark.sql.inMemoryColumnStorage.batchSize:默认值为1000,代表列式缓存时每个批处理的大小。如果将这个值调得过大,可能会产生内存溢出(Out Of Memory,OOM)的异常。因此,在设置这个的参数时要注意实际的内存大小。

       上面两个参数在之前的多个版本都能使用。
       下面几个参数可能在之后的新版本中保留或者删除。

  • spark.sql.files.maxPartitionBytes:默认值是134217728(128MB),这个参数代表partition的最大数。
  • spark.sql.files.openCostInBytes:默认值是4194304(4MB),这个参数代表小于4MB的文件会被合并到一个partition中。
  • spark.sql.broadcastTimeout:默认值是300,广播的超时时间,以秒为单位。
  • spark.sql.autoBroadcastJoinThreshold:默认值是10485760(10MB),表示大表.join(小表)。读者可以根据需要广播的小表调整参数的大小。当使用连接操作时,会自动将小于阈值的表广播给所有工作节点。利用好这个属性可以降低数据传输的网络开销。当这个属性的值被设为-1时,关闭广播。
  • spark.sql.shuffle.partitions:默认值是200,这个参数代表着执行连接操作或聚合操作时数据分区的数目(由于计算是以partition为单位进行的,因此称之为并行度,这里需要读者稍微注意一下)。读者根据实际情况调大或调小,找到合适值就好。该参数在一定程度上能减少数据倾斜。

上面这些都是针对Spark SQL的属性。接下来讲的是针对整个Spark的属性。

  • spark.default.parallelism:这个是Spark的默认并行度(就是默认的分区数目)。对于不同的环境,默认配置不一样:
  • 对于local模式来说,这个属性的值是机器上的核心数。
  • 对于Mesos的细粒度模式来说,这个属性的值是8。
  • 对于其他的资源管理器,比如YARN,对应的值是所有执行节点的核心数,最低是2。
    通常,每个CPU核分配2~3个Task。对于有超线程技术的CPU,还是要以核心数为主,而不是核心数×2。
  • park.executor.core和spark.executor.memory:这两个属性用于修改每个Executor使用的核心数以及内存大小。
  • spark.dirver.core和spark.dirver.memory:这两个属性用于修改每个驱动进程使用的核心数以及内存大小。
  • spark.executor.instances:这个属性用于设置启动的Executor数目。

提示: http:// :4040这个页面给我们提供了很多非常实用的信息,希望读者能花一些时间熟悉这个Web UI。在该页面的Environment栏中可以查看已经生效的属性,未显示在上面的属性则认为是使用了默认值。由于配置属性是随着版本的发行而经常变动的,因此这里就不详细讲述了,读者可根据自己使用的Spark版本查阅相应的官方文档。

下面举一个使用Spark on YARN的例子。
在该集群上有6台主机运行着NodeManager,每台主机有16个核心和64GB的内存。

  • yarn.nodemanager.resource.memory-mb设置为63×1024MB=64512MB=63GB。
  • yarn.nodemanager.resource.cpu-vcores设置为15。
    为什么不设置为64×1024MB=65536MB的内存和16个核心呢?因为系统运行需要内存,Hadoop的守护进程也需要内存,所以要留1GB和1个核心给它们。然后读者可能会使用以下 配置:
--num-executors 6  --executor-cores 15  --executor-memory 63G

       这仍然是不妥当的,因为我们还需要考虑Executor的内存开销。63GB分配给Executor使用,再加上Executor本身的内存开销,就超过了分配给NodeManager的63GB内存。
       除此之外,我们还需要考虑ApplicationMaster的CPU使用。ApplicationMaster本身需要占用一个核心,剩下的就不够15个核心了,也就是说分配不到15个核心给Executor。
此外,给一个Executor分配15个核心还会导致HDFS的I/O吞吐量变得很差。
下面是改良之后的参考配置:

--num-executors 17 --executor-cores 5 --executor-memory 19G

       使用上面这种配置,5台主机上每台都有3个Executor。最后一台主机上只有两个Executor,这是因为这台主机上还运行着ApplicationMaster。
       关于内存的计算如下:

63GB/3=21.21GB
21.21GB*0.07=1.47GB
21.21GB-1.47GB≈19GB

       上面的0.07是怎么来的呢?这与spark.yarn.executor.memoryOverhead这个参数有关:在Spark 2.2.1版本中,该属性的取值默认是0.10*executorMemory,低于384MB按384MB计算。对于0.10,这里可以取其他值,比如0.06~0.10,上面的0.07的作用与0.10类似。

13.3.4 合理使用广播

       对于比较大的变量,我们可以将它广播到每一个节点中,以节省网络通信的开销。在不广播的情况下,每个Task有一个数据的副本,在广播之后每个Executor保留一份数据的副本。因为广播之后减少了数据副本的数量,所以在减少网络传输开销的同时也相应地节省了一些资源。
       广播的使用方式如下:

val broadcastVar = sc.broadcast(Array(1, 2, 3))  // 这样我们就把Array(1,2,3)广播到各个节点中去了
broadcastVar.value // 这样就可以调用Array(1,2,3)

13.3.5 尽量避免使用Shuffle算子

       Shuffle操作涉及磁盘的I/O操作、数据的序列化和网络的I/O。某些Shuffle操作会消耗大量的内存,它会把相同的key发到一个节点中,进行连接或者聚合操作。当具有相同key的数据量特别大时,内存有可能溢出,于是将数据写到磁盘上,引发I/O操作,导致性能急剧下降。
       典型地使用了Shuffle操作的有repartition、coalesce、groupByKey、reduceByKey、coGroup、Join等。如果因为硬性需求必须使用带Shuffle的操作,那么尽量使用在Map端就聚合一次的方法,比如用reduceByKey或者aggregateByKey替代groupByKey。
       那么,什么是在Map端聚合呢?下面来看两幅图。
       图13-3所示是rdd.reduceByKey(_ + _)的图。
在这里插入图片描述
图13-4所示是rdd.groupByKey().map(t => (t._1 , t._2.sum))的图。
在这里插入图片描述
       相信读者看了这两幅图就大概能明白为什么推荐用reduceByKey替代groupByKey了。groupByKey是将key相同的数据发送到一个节点中,然后在那个节点上进行值的相加操作。如果key相同的数据过多,就会增加网络的开销和内存的开销。相比之下,reduceByKey在Map端对key相同的数据进行了值相加的操作,得出一个中间结果,然后将中间结果中key相同的数据发送到同一个节点中,再将它们的值相加得出最终的结果。如此一来,就减少了需要通过网络传输的数据量,同时也节省了Reduce端内存的开销。
       下面用一个例子来说明一下。
       例如以下代码:

rdd.map(kv => (kv._1, new Set[String]() + kv._2)).reduceByKey(_ ++ _)

       这行代码对每一条记录进行处理时,都会创建一个Set对象。这和前面提到的六大开发原则中的第一条相悖,要避免重复创建不必要的对象。另外,这里还使用了reduceByKey,前文提到了在使用含有Shuffle操作的方法时一定要谨慎,一定要理解为什么用它,并且考虑有没有更好的方法。
       改良之后的代码如下:

val empty = new collection.mutable.Set[String]()
rdd.aggregateByKey(empty)((set, v) => set += v,(set1, set2) => set1 ++= set2)

       在这段代码中,我们用aggregateByKey代替了reduceByKey,最终的效果是一样的,但却巧妙地减少了许多不必要对象的创建,大大提高了执行效率。

13.3.6 使用map-side预聚合的Shuffle操作

       如果因为业务需要,一定要使用Shuffle操作,而无法用Map类的算子来替代,那么尽量使用可以进行map-side预聚合的算子。所谓的map-side预聚合,就是在每个节点本地对相同的key进行一次聚合操作,类似于MapReduce中的本地combiner。经过map-side预聚合之后,每个节点本地只会有一条相同的key,因为多条相同的key都被聚合起来了。其他节点在拉取所有节点上的相同的key时,就会大大减少需要拉取的数据量,从而减少磁盘I/O以及网络传输的开销。
       通常来说,在可能的情况下,建议使用reduceByKey或者aggregateByKey算子来替代groupByKey算子。因为reduceByKey和aggregateByKey算子都会使用用户自定义函数,对每个节点本地的相同key进行预聚合;而groupByKey算子是不会进行预聚合的,全量的数据会在集群的各个节点之间分发和传输,性能相对来说比较差。例如,图13-5所示就是典型的例子,分别基于reduceByKey和groupByKey进行单词计数。图的上半部分是groupByKey的原理图,可以看到没有进行任何本地聚合时,所有数据都会在集群节点之间传输;图的下半部分是reduceByKey的原理图,可以看到每个节点本地的相同key数据都进行了预聚合,然后才传输到其他节点上进行全局聚合。
在这里插入图片描述

13.3.7 使用高性能算子

       除了Shuffle相关的算子有优化原则之外,其他的算子也有着相应的优化原则。
       使用mapPartitions替代普通Map。mapPartitions类的算子,一次函数调用会处理一个分区中的所有数据,而不是只处理一条数据,性能相对来说会高一些。但是,有的时候使用mapPartitions会出现内存溢出的问题。因为单次函数调用就要处理掉一个分区中的所有数据,如果内存不够,垃圾回收时无法回收太多对象,很可能出现内存溢出异常。因此,使用这类操作时要慎重。
       使用foreachPartitions替代foreach的原理类似于“使用mapPartitions替代Map”,也是一次函数调用处理一个分区中的所有数据,而不是只处理一条数据。在实践中发现,foreachPartitions类的算子对性能的提升很有帮助。例如在foreach函数中,将RDD中所有数据写入MySQL,如果是普通的foreach算子,就会一条数据一条数据地写,每次函数调用可能都会创建一个数据库连接,这样势必会频繁地创建和销毁数据库连接,性能非常低下。但是,如果使用foreachPartitions算子一次性处理一个分区的所有数据,那么对于每个分区,只要创建一个数据库连接即可,然后执行批量插入操作,此时性能比较高。在实践中发现,对于将1万条左右的数据写入MySQL,使用foreachPartitions算子的性能可以提升30%以上。
       通常对一个RDD执行filter算子,过滤掉RDD中的较多数据后(比如30%以上的数据),建议使用coalesce算子,手动减少RDD的分区数量,将RDD中的数据压缩到更少的分区中去。因为使用filter之后,RDD的每个分区中都会有很多数据被过滤掉,此时如果照常进行后续的计算,其实每个Task处理的分区中的数据量并不是很多,有一点浪费资源,而且此时处理的Task越多,速度可能反而越慢。因此,用coalesce减少分区数量,将RDD中的数据压缩到更少的分区之后,只要使用更少的Task即可处理完所有的分区。在某些场景下,这对于性能的提升会有一定的帮助。
       使用repartitionAndSortWithinPartitions替代repartition与sort类操作repartitionAndSortWithinPartitions是Spark官网推荐的一个做法。官方建议,如果需要在重分区之后进行排序,建议直接使用repartitionAndSortWithinPartitions算子。因为该算子可以一边进行重分区的Shuffle操作,一边进行排序。Shuffle与Sort两个操作同时进行,比先进行Shuffle操作再进行Sort操作性能要高。

13.3.8 尽量在一次调用中处理一个分区的数据

       mapPartitions、foreachPartitions都是在一次调用中处理一个分区的数据,所以用mapPartitions替代map、用foreachPartitions替代foreach能提高性能。但是使用时需要注意内存,因为一次处理一个分区,如果分区比较大,当内存不够时就会出现内存溢出异常。
       map和mapPartitions在使用上是有区别的,下面举例说明:

val rdd = sc.parallelize(1 to 9, 3)// 分3个分区的RDD
def mapFunc(num:Int):Int = {
           var result = num*num
           result
      }
    def mapPartitionsFunc ( iter : Iterator [Int] ) : Iterator [Int] = {
          var result = for (num <- iter ) yield num*num
          result
    }
rdd.map(mapFunc)
rdd.mapPartitions(mapPartitionsFunc)

       在上面这段代码中,mapFunc被执行了10次,而mapPartitionsFunc被执行了3次。还有就是mapFunc和mapPartitionsFunc传入的参数不一样,这个需要在使用时注意一下。
       之前提到的用mapPartitions替代map、用foreachPartitions替代foreach会提高性能,这是为什么呢?假如我们在上述代码的mapFunc和mapPartitionsFunc中创建相同的对象或者创建数据库连接,在mapFunc中会被创建10次,而在mapPartitionsFunc中只创建3次,这就大大地减少了内存开销。

13.3.9 对数据进行序列化

       对数据进行序列化可以使数据更紧凑、更小,以此减少网络的传输开销,但是会使访问对象的时间变长,因为数据进行反序列化之后才能使用。
       就现在来说,Spark的默认数据序列化方式是调用Java的ObjectOutputStream框架。如果读者想提高序列化的效率,就可以使用Kryo。
       示例代码如下:

// conf是SparkConf的实例
conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer" )
// 下面是序列化自定义的对象
conf.registerKryoClasses(Array(classOf[MyClass1], classOf[MyClass2]))
val sc = new SparkContext(conf)

如果序列化对象太大,那么可以设置spark.kryoserializer.buffer来进行调整。
关于Kryo的更多信息可以在https:// github.com/EsotericSoftware/kryo中找到。

13.4 Spark数据倾斜调优

       相信很多读者都听说过数据倾斜,那么数据倾斜到底是什么呢?我们知道,在进行Shuffle操作时会将各个节点上key相同的数据传输到同一节点以进行下一步操作。如果某个key或某几个key下的数据量特别大,远远大于其他key的数据,这时就会出现一个现象:大部分Task很快就完成,剩下几个Task运行特别缓慢,甚至有时还会因为某个Task下相同key的数据量过大而造成内存溢出。这就是发生了数据倾斜。
       兵来将挡,水来土掩。既然是数据发生了倾斜,那么主要的解决思路就是想办法让它不 倾斜。

13.4.1 调整分区数目

       前提条件是Task上可以分配多个key的数据。
       在发生数据倾斜时,某个Task上需要处理的数据过多,我们可以调整并行度,使原本分配给一个Task的多个key分配给多个Task,这样需要Task处理的key的数目就会减少,于是Task上的数据量也就减少了。
       在13.3节中,我们提到过一个配置属性spark.sql.shuffle.partitions,一般是把它的值适当地调大。这个方法只能缓解数据倾斜,没有从根源上解决问题。但是这个方法比较简单,推荐优先尝试使用。
       另外,在进行多次操作之后,会有很多小任务产生,这时可以用coalesce来减少分区数。但分区数不是越少越好,当数据是几个特别大的,并且不可分的文件时,如果分区过少,就不能充分使用CPU的所有核心。这种情况下,就需要主动(触发一次Shuffle)进行重分区,增加分区的数量,以提高并行度。
       有很多方法都提供了可以调整分区数目的参数,例如下面这个。

val rdd2 = rdd1.reduceByKey(_ + _, numPartitions = X)

       分区数目应该怎么设置呢?一般来说需要一点一点地尝试,比如按父分区数×1.5这样一点一点地往上调,直到性能足够好时停止增加。
每个任务可用的内存是:

(spark.executor.memory * spark.shuffle.memoryFraction * spark.shuffle.safetyFraction)/spark.executor.cores

查看分区数的代码如下:

rdd.partitions().size()

13.4.2 去除多余的数据

首先,查看每个key的数据量,代码如下:

pairs.sample(false,0.1).countByKey().foreach(println())

如果发现导致数据倾斜的部分key对最后的结果没有影响,就过滤掉这些数据,从而避免数据倾斜的发生。

13.4.3 使用广播将Reduce Join转换为Map Join

(1)调整spark.sql.autoBroadcastJoinThreshold的大小,使其大于需要广播的小表,这样就会自动广播小表。
(2)使用Broadcast广播小表。
这样可以避免发生Shuffle操作,示例如下:

// 大表字段:id  class  score
// 小表字段:id  name
// 输出结果表字段:name  class  score
// 下面的代码在Map端执行join,不经历Shuffle和Reduce,执行效率比较高
var broadcastTable = sc.broadcast(aSmallTable)// aSmallTable 是Map组成的RDD
var result = bigTable.mapPartition( iter=>{
	var smallTable =  broadcastTable.value
	var arrayBuffer  = ArrayBuffer[(String,String,String)]() 
	iter.foreach{case(id,class,score)=>{
		if(smallTable.contain(id)){
			 arrayBuffer += ((smallTable.getOrElse(id,""),class,score))
		}
	}}
	 arrayBuffer.iterator
})

13.4.4 将key进行拆分,大数据转换为小数据

       回顾一下数据倾斜的原因——单个key或某几个key的数据过多。既然数据过多,就想办法减少单个key的数据量。我们可以给key加上前缀,强行让它们不同。通过这种方式让本来应该到同一个Task的数据分散到不同的Task上,以此来解决数据倾斜的问题。
       需要注意的是,对聚合操作的key的拆分和对Join操作的key的拆分是不一样的。下面用示意图来说明它们的原理以及不同之处。
不拆解key时的聚合操作如图13-6所示。
在这里插入图片描述
       对于聚合操作的key的拆解如图13-7所示。
       图13-7展示的是对相同key随机加了前缀,于是将key拆分,使它们分布在不同的Task上,然后逐步聚合。这样就能比较有效地化解数据倾斜的问题。
       对于Join操作的key的拆解如图13-8所示。
       这里的原理其实和上面聚合的原理差不多,只是需要注意:一张表加n种前缀,另一张表则翻n倍,这样才能正常地连接。但这样会使内存资源的消耗翻倍,因此使用时要仔细斟酌。
在这里插入图片描述
       其实解决数据倾斜这个问题有多种技巧,这里给读者列出的是容易上手的几种。剩下的需要读者根据实际的环境和需求进行优化。

13.4.5 数据倾斜定位和解决

1. 数据倾斜
       Spark中的数据倾斜问题主要指的是Shuffle过程中出现的数据倾斜问题,是由于不同的key对应的数据量不同导致不同Task所处理的数据量不同的问题。
例如,Reduce点一共要处理100万条数据,第一个和第二个Task分别被分配到了1万条数据,计算5分钟内完成,第三个Task分配到了98万条数据,此时第三个Task可能需要10个小时完成,这使得整个Spark作业需要10个小时才能运行完成,这就是数据倾斜所带来的后果。

2. 数据倾斜有两大直接致命后果
(1)数据倾斜直接会导致一种情况:内存溢出(OOM)。
(2)运行速度慢。
注意,要区分开数据倾斜与数据量过量这两种情况。数据倾斜是指少数Task被分配了绝大多数的数据,因此少数Task运行缓慢;数据过量是指所有Task被分配的数据量都很多,相差不大,所有Task都运行缓慢。

       数据倾斜的表现:
       (1)Spark作业的大部分Task都执行得很迅速,只有有限的几个Task执行得非常慢,此时可能会出现数据倾斜问题,作业可以运行,但是运行得非常缓慢。
       (2)Spark作业的大部分Task都执行得很迅速,但是有的Task在运行过程中会突然报出OOM错误,反复执行几次都在某一个Task报出OOM错误,此时可能出现了数据倾斜问题,作业无法正常运行。

       定位数据倾斜问题:
       (1)查阅代码中的Shuffle算子,例如reduceByKey、countByKey、groupByKey、Join等算子,根据代码逻辑判断此处是否会出现数据倾斜问题。
       (2)查看Spark作业的Log文件,Log文件对于错误的记录会精确到代码的某一行,可以根据异常定位到的代码位置来明确错误发生在第几个Stage,对应的Shuffle算子是哪一个。

3. 解决方案一:聚合元数据
       1)避免shuGle过程
       绝大多数情况下,Spark作业的数据来源都是Hive表,这些Hive表基本都是经过ETL之后的昨天的数据。为了避免数据倾斜问题,我们可以考虑避免Shuffle过程,如果避免了Shuffle过程,那么从根本上就消除了发生数据倾斜问题的可能。
       如果Spark作业的数据来源于Hive表,那么可以先在Hive表中对数据进行聚合。例如,按照key进行分组,将同一key对应的所有value用一种特殊的格式拼接到一个字符串中,这样一个key就只有一条数据了;之后,在对一个key的所有value进行处理时,只需要进行Map操作即可,无须再进行任何的Shuffle操作。通过上述方式就避免了执行Shuffle操作,也就不可能发生任何数据倾斜问题。
       2)缩小key粒度(增大数据倾斜可能性,降低每个Task的数据量)
       key的数量增加,可能使数据倾斜更严重。

       3)增大key粒度(减小数据倾斜可能性,增大每个Task的数据量)
       如果没有办法对每个key聚合出一条数据,在特定场景下,可以考虑扩大key的聚合粒度。
例如,有10万条用户数据,当前key的粒度是(省,城市,区,日期)。现在我们考虑扩大粒度,将key的粒度扩大为(省,城市,日期)。这样的话,key的数量会减少,key之间的数据量差异也有可能会减少,由此可以减轻数据倾斜的现象和问题。此方法只针对特定类型的数据有效,当应用场景不适宜时,会加重数据倾斜。

4. 解决方案二:过滤导致倾斜的key
如果在Spark作业中允许丢弃某些数据,那么可以考虑将可能导致数据倾斜的key进行过滤,滤除可能导致数据倾斜的key对应的数据。这样,在Spark作业中就不会发生数据倾斜了。

5. 解决方案三:提高shuGle操作中的reduce并行度
       当方案一和方案二对于数据倾斜的处理没有很好的效果时,可以考虑提高Shuffle过程中的Reduce端并行度。Reduce端并行度的提高也就增加了Reduce端Task的数量,那么每个Task分配到的数据量就会相应减少,由此缓解数据倾斜问题。
       1)Reduce端并行度的设置
       在大部分的Shuffle算子中,都可以传入一个并行度的设置参数,比如reduceByKey(500),这个参数会决定Shuffle过程中Reduce端的并行度。在进行Shuffle操作时,就会对应创建指定数量的Reduce Task。对于Spark SQL中的Shuffle类语句,比如group by、Join等,需要设置一个参数,即spark.sql.shuffle.partitions,该参数代表了Shuffle Read Task的并行度,该值默认是200,对于很多场景来说有点过小。增加Shuffle Read Task的数量可以让原本分配给一个Task的多个key分配给多个Task,从而让每个Task处理比原来更少的数据。
举例来说,如果原本有5个key,每个key对应10条数据,这5个key都是分配给一个Task的,那么这个Task就要处理50条数据。而增加了Shuffle Read Task后,每个Task就分配到一个key,即每个Task就处理10条数据,那么自然每个Task的执行时间都会变短了。
       2)Reduce端并行度设置存在的缺陷
       提高Reduce端并行度并没有从根本上改变数据倾斜的本质和问题(方案一和方案二从根本上避免了数据倾斜的发生),只是尽可能地缓解和减轻Shuffle Reduce Task的数据压力,以及数据倾斜问题。这种方案适用于有较多key对应的数据量都比较大的情况,该方案通常无法彻底解决数据倾斜问题,因为如果出现一些极端情况,比如某个key对应的数据量有100万条,那么无论你的Task数量增加到多少,这个对应100万条数据的key肯定还是会被分配到一个Task中处理, 因此注定还是会发生数据倾斜问题。因此,这种方案只能说是在发现数据倾斜问题时尝试使用的第一种手段,尝试用最简单的方法缓解数据倾斜而已,也可以和其他方案结合起来使用。
在理想情况下,Reduce端的并行度提升后,会在一定程度上减轻数据倾斜问题,甚至基本消除数据倾斜问题。而且,在一些情况下,只会让原来由于数据倾斜而运行缓慢的Task运行速度稍有提升,或者避免某些Task的OOM问题,但是仍然运行缓慢。此时要及时放弃方案三,开始尝试后面的方案。

6. 解决方案四:使用随机key实现双重聚合
       当使用了类似于groupByKey、reduceByKey这样的算子时,可以考虑使用随机key实现双重聚合。
       首先,通过Map算子给每个数据的key添加随机数前缀,将key打散,将原先一样的key变成不一样的key;然后进行第一次聚合,这样就可以将原本被一个Task处理的数据分散到多个Task上进行局部聚合。
       随后,去除掉每个key的前缀,再次进行聚合。此方法对于由groupByKey、reduceByKey这类算子造成的数据倾斜问题有比较好的效果,它仅仅适用于聚合类的Shuffle操作,适用范围相对较窄。如果是Join类的Shuffle操作,还得用其他的解决方案,此方法也是在前几种方案中没有比较好的效果时可以尝试的解决方法。

7. 解决方案五:将Reduce Join转换为Map Join
       正常情况下,Join操作都会执行Shuffle过程,并且执行的是Reduce Join,也就是先将所有相同的key和对应的value汇聚到一个Reduce Task中,然后进行Join。普通的Join是会进行Shuffle过程的,而一旦Shuffle,就相当于会将相同key的数据拉取到一个Shuffle Read Task中再进行Join,此时就是Reduce Join。但是如果一个RDD比较小,则可以采用广播小RDD全量数据+Map算子来实现与Join同样的效果,也就是Map Join,此时就不会进行Shuffle操作,也就不会发生数据倾斜问题。
       注意,RDD是不能进行广播的,只能将RDD内部的数据通过collect算子拉取到Driver内存中,然后进行广播。
1)核心思路
       不使用Join算子进行连接操作,而是使用Broadcast变量与Map类算子实现Join操作,进而完全规避掉Shuffle类的操作,彻底避免数据倾斜问题的发生和出现。将较小RDD中的数据直接通过collect算子拉取到Driver端的内存中,然后对其创建一个Broadcast变量;接着对另一个RDD执行Map类算子,在算子函数内,从Broadcast变量中获取较小RDD的全量数据,与当前RDD的每一条数据按照连接key进行比对,如果连接key相同的话,那么就将两个RDD的数据用你需要的方式连接起来。
       根据上述思路,根本不会发生Shuffle操作,从根本上杜绝了Join操作可能导致的数据倾斜问题。当Join操作有数据倾斜问题并且其中一个RDD的数据量较小时,可以优先考虑这种方式,效果非常好。
2)不使用场景分析
       由于Spark的广播变量是在每个Executor中保存一个副本,如果两个RDD数据量都比较大,那么如果将一个数据量比较大的RDD做成广播变量,那么很有可能造成内存溢出。

8. 解决方案六:sample采样对倾斜key单独进行Join
       在Spark中,如果某个RDD只有一个key,那么在Shuffle过程中会默认将此key对应的数据打散,由不同的Reduce端Task进行处理。当由单个key导致数据倾斜时,可以将发生数据倾斜的key单独提取出来,组成一个RDD,然后用这个原本会导致倾斜的key组成的RDD跟其他RDD单独Join。此时,根据Spark的运行机制,此RDD中的数据会在Shuffle阶段被分散到多个Task中进行Join操作。
1)适用场景分析
       对于RDD中的数据,可以将其转换为一个中间表,或者直接使用countByKey()的方式,看一下这个RDD中各个key对应的数据量,此时如果你发现整个RDD只有一个key的数据量特别多,那么就可以考虑使用这种方法。
当数据量非常大时,可以考虑使用sample采样获取10%的数据,然后分析这10%的数据中哪个key可能会导致数据倾斜,然后将这个key对应的数据单独提取出来。
2)不适用场景分析
       如果一个RDD中导致数据倾斜的key很多,那么此方案不适用。

9. 解决方案七:使用随机数以及扩容进行Join
       如果在进行Join操作时,RDD中有大量的key导致数据倾斜,那么分拆key也没什么意义,此时只能使用最后一种方案来解决问题。对于Join操作,我们可以考虑对其中一个RDD数据进行扩容,另一个RDD进行稀释后再Join。我们会将原先一样的key通过附加随机前缀变成不一样的key,然后就可以将这些处理后的“不同key”分散到多个Task中处理,而不是让一个Task处理大量的相同key。这种方案是针对有大量倾斜key的情况,没法将部分key拆分出来单独进行处理,需要对整个RDD进行数据扩容,对内存资源要求很高。
1)核心思想
       选择一个RDD,使用flatMap进行扩容,对每条数据的key添加数值前缀(1~N的数值),将一条数据映射为多条数据。
       (扩容)选择另一个RDD,进行Map映射操作,每条数据的key都打上一个随机数作为前缀(1~N的随机数)。
       (稀释)将两个处理后的RDD进行Join操作。
2)局限性
       如果两个RDD都很大,那么将RDD进行N倍的扩容显然行不通,使用扩容的方式只能缓解数据倾斜,不能彻底解决数据倾斜问题。
3)使用方案七对方案六进一步优化分析
       当RDD中有几个key导致数据倾斜时,方案六不再适用,而方案七又非常消耗资源,此时可以引入方案七的思想来完善方案六。
       (1)对包含少数几个数据量过大的key的那个RDD,通过sample算子采样出一份样本,然后统计一下每个key的数量,计算出数据量最大的是哪几个key。
       (2)然后将这几个key对应的数据从原来的RDD中拆分出来,形成一个单独的RDD,并给每个key都打上n以内的随机数作为前缀,而不会导致倾斜的大部分key形成另一个RDD。
       (3)接着将需要Join的另一个RDD,也过滤出那几个倾斜key对应的数据并形成一个单独的RDD,将每条数据膨胀成n条数据,这n条数据都按顺序附加一个0~n的前缀,不会导致倾斜的大部分key也形成另一个RDD。
       (4)再将附加了随机前缀的独立RDD与另一个膨胀n倍的独立RDD进行Join,此时就可以将原先相同的key打散成n份,分散到多个Task中进行Join。
       (5)而另外两个普通的RDD照常进行Join即可。
       (6)最后将两次Join的结果使用union算子合并起来,就是最终的Join结果。

13.5 本章小结

       本章全面探讨了Spark性能调优的策略和方法。首先,从常规性能调优入手,介绍了调整Spark配置参数、优化资源分配等基础手段。接着,深入探讨了Spark开发优化原则,强调了编写高效Spark程序的重要性,包括避免不必要的数据传输、优化数据处理逻辑等。在Spark调优方法部分,详细讲解了如何通过监控和分析Spark作业的运行情况,找出性能瓶颈并进行针对性优化。最后,本章针对Spark数据倾斜这一常见问题,提供了多种调优策略,比如使用Salting技术、调整分区策略等,以缓解数据倾斜带来的性能影响。通过本章的学习,读者将能够掌握Spark性能调优的关键技巧,为提升Spark作业的运行效率和性能提供有力支持。

本书其他章节

  1. Spark大数据开发与应用案例(视频教学版)(一)–文前
  2. Spark大数据开发与应用案例(视频教学版)(二)–第一章上
  3. Spark大数据开发与应用案例(视频教学版)(三)–第一章下
  4. Spark大数据开发与应用案例(视频教学版)(四)–第二章上
  5. Spark大数据开发与应用案例(视频教学版)(五)–第二章下
  6. Spark大数据开发与应用案例(视频教学版)(六)–第三章上
  7. Spark大数据开发与应用案例(视频教学版)(七)–第三章下
  8. Spark大数据开发与应用案例(视频教学版)(八)–第四章上
  9. Spark大数据开发与应用案例(视频教学版)(九)–第四章下
  10. Spark大数据开发与应用案例(视频教学版)(十)–第五章
  11. Spark大数据开发与应用案例(视频教学版)(十一)–第六章
  12. Spark大数据开发与应用案例(视频教学版)(十二)–第七章
  13. Spark大数据开发与应用案例(视频教学版)(十三)–第八章
  14. Spark大数据开发与应用案例(视频教学版)(十四)–第九章
  15. Spark大数据开发与应用案例(视频教学版)(十五)–第十章上
  16. Spark大数据开发与应用案例(视频教学版)(十六)–第十章下
  17. Spark大数据开发与应用案例(视频教学版)(十七)–第十一章
  18. Spark大数据开发与应用案例(视频教学版)(十八)–第十二章
  19. Spark大数据开发与应用案例(视频教学版)(十九)–第十三章
  20. Spark大数据开发与应用案例(视频教学版)(二十)–第十四章
  21. Spark大数据开发与应用案例(视频教学版)(二十一)–第十五章

点赞+收藏+关注

在这里插入图片描述

更多推荐