暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Spark 内核解析

BigData Scholar 2021-06-29
606

1.1 概述

是一种基于内存的快速、通用、可扩展的大数据分析引擎

1.2 内置模块

1.3 特点

1.3.1 快

  • 基于内存计算,中间结果存于内存中
  • 高效DAG执行引擎

1.3.2 易用

  • 支持scala、java、python的API,支持scala和python的Shell

1.3.3 通用

  • 批处理、交互式查询(sql)、实时流处理(streming)、机器学习(MLlib)、图计算(GraphX)

13.4 兼容性

  • 与其他开源产品融合

    1)资源管理与调度:yarn、mesos

    2)数据处理:hdfs、hbase、hive等

1.4 重要角色

1.4.1 Driver驱动器

  • 概念

    执行开发程序中的main方法的进程。负责执行创建SparkContext、创建RDD,以及进行RDD的转化操作和行动操作的代码的进程

  • 主要职责

    1)把用户程序转为作业(JOB)

    2)跟踪Executor的运行状况

    3)为执行器节点调度任务

    4)UI展示应用运行状况

1.4.2 Executor执行器

  • 概念

    工作进程,负责在 Spark 作业中运行任务,任务间相互独立。Spark 应用启动时,Executor节点被同时启动,并且始终伴随着整个 Spark 应用的生命周期而存在。如果有Executor节点发生了故障或崩溃,Spark 应用也可以继续执行,会将出错节点上的任务调度到其他Executor节点上继续运行

  • 主要职责

    1)负责运行组成 Spark 应用的任务,并将结果返回给驱动器进程;

    2)通过自身的块管理器(Block Manager)为用户程序中要求缓存的RDD提供内存式存储。RDD是直接缓存在Executor进程内的,因此任务可以在运行时充分利用缓存数据加速运算。

1.5 经典实例:wordcount

1.5.1 代码

def main(args: Array[String]): Unit = {

//1.创建SparkConf并设置App名称
    val conf = new SparkConf().setAppName("WordCount")
//2.创建SparkContext,sc是提交Spark App的入口
    val sc = new SparkContext(conf)
    //3.使用sc创建RDD并执行相应的transformation和action
    sc.textFile(args(0)).flatMap(_.split(" ")).map((_, 1)).reduceByKey(_+_, 1).sortBy(_._2, false).saveAsTextFile(args(1))
//4.关闭连接
    sc.stop()
  }

15.2 shell调用

bin/spark-submit \
--class <main-class>
--master <master-url> \
--deploy-mode <deploy-mode> \
--conf <key>=<value> \
... # other options
<application-jar> \
[application-arguments]\
--executor-memory 1G\
--total-executor-cores 2

  • --master 指定Master的地址,默认为Local
  • --class: 你的应用的启动类 (如 org.apache.spark.examples.SparkPi)
  • --deploy-mode: 是否发布你的驱动到worker节点(cluster) 或者作为一个本地客户端 (client) (default: client)*
  • --conf: 任意的Spark配置属性, 格式key=value. 如果值包含空格,可以加引号“key=value”
  • application-jar: 打包好的应用jar,包含依赖. 这个URL在集群中全局可见。比如hdfs:// 共享存储系统, 如果是 file:// path, 那么所有的节点的path都包含同样的jar
  • application-arguments: 传给main()方法的参数
  • --executor-memory 1G 指定每个executor可用内存为1G
  • --total-executor-cores 2 指定每个executor使用的cup核数为2个

1.6 通用运行流程

1.7 运行模式&调度机制

1.7.1 local 模式

    运行在一台计算机上

1.7.2 Standalone模式

    由master+slave构成的spark集群

  • 任务提交流程

1.7.3 yarn模式[重点]

    客户端直接连接yarn,不需要额外构建spark集群,有yarn-client和yarn-cluster两种模式,主要区别在于Driver程序的运行节点

1.7.3.1 yarn-client

    Driver程序运行在客户端,适用于交互、调试,希望立即看到app的输出

  • 任务提交流程:

1.7.3.2 yarn-cluster[重点]

    Driver程序运行在由ResourceManager启动的APPMaster上,适用于生产环境

1.7.3.2.1 任务提交流程

1.7.3.2.2 任务调度
  • Spark应用程序包括Job、Stage、Task三个概念

    1)Job是以Action方法为界,遇到一个Action方法则触发一个Job;

    2)Stage是Job的子集,以RDD宽依赖(即Shuffle)为界,遇到Shuffle做一次划分;

    3)Task是Stage的子集,以并行度(分区数)来衡量,分区数是多少,则有多少个task

  • 概述

    当Driver起来后,Driver会根据用户程序逻辑准备任务,并根据Executor资源情况逐步分发任务。Spark RDD通过其Transactions操作,形成了RDD血缘关系图,即DAG,最后通过Action的调用,触发Job并调度执行。DAGScheduler负责Stage级的调度,主要是将job切分成若干Stages,并将每个Stage打包成TaskSet交给TaskScheduler调度。TaskScheduler负责Task级的调度,将DAGScheduler给过来的TaskSet按照指定的调度策略分发到Executor上执行,调度过程中SchedulerBackend负责提供可用资源,其中SchedulerBackend有多种实现,分别对接不同的资源管理系统

  • Stage级调度

    Job由最终的RDD和Action方法封装而成,SparkContext将Job交给DAGScheduler提交,它会根据RDD的血缘关系构成的DAG进行切分,将一个Job划分为若干Stages,具体划分策略是,由最终的RDD不断通过依赖回溯判断父依赖是否是宽依赖,即以Shuffle为界,划分Stage,窄依赖的RDD之间被划分到同一个Stage中,可以进行pipeline式的计算。划分的Stages分两类,一类叫做ResultStage,为DAG最下游的Stage,由Action方法决定,另一类叫做ShuffleMapStage,为下游Stage准备数据

  • Task级调度

    Spark Task的调度是由TaskScheduler来完成,DAGScheduler将Stage打包到TaskSet交给TaskScheduler,TaskScheduler会将TaskSet封装为TaskSetManager加入到调度队列中,TaskSetManager负责监控管理同一个Stage中的Tasks,TaskScheduler就是以TaskSetManager为单元来调度任务

1)调度策略

(1)FIFO

    将TaskSetManager按照先来先到的方式入队,出队时直接拿出最先进队的TaskSetManager

(2)Fair

    FAIR模式中有一个rootPool和多个子Pool,各个子Pool中存储着所有待分配的TaskSetMagager。需要先对子Pool进行排序,再对子Pool里面的TaskSetMagager进行排序(Pool和TaskSetMagager都继承了Schedulable特质,因此使用相同的排序算法)。排序过程的比较是基于Fair-share来比较的,每个要排序的对象包含三个属性: runningTasks值(正在运行的Task数)、minShare值(最小资源分配值)、weight值(可以分配的资源比例)。

     ① runningTasks比minShare小的先执行

     ② minShare使用率低的先执行

     ③ 权重使用率低的先执行

2)本地化调度

     ① PROCESS_LOCAL 进程本地化,task和数据在同一个Executor中,性能最好。

     ② NODE_LOCAL 节点本地化,task和数据在同一个节点中,但是task和数据不在同一个Executor中,数据需要在进程间进行传输。

     ③ RACK_LOCAL 机架本地化,task和数据在同一个机架的两个节点上,数据需要通过网络在节点之间进行传输。

     ④ NO_PREF 对于task来说,从哪里获取都一样,没有好坏之分。

     ⑤ ANY task和数据可以在集群的任何地方,而且不在一个机架中,性能最差。

3)失败重试与黑名单机制

    对于失败的Task,会记录它失败的次数,如果失败次数还没有超过最大重试次数,那么就把它放回待调度的Task池子中,否则整个Application失败。同时记录它上一次失败所在的Executor Id和Host,下次再调度这个Task时,会使用黑名单机制,避免它被调度到上一次失败的节点上,起到一定的容错作用

1.8 Shuffle

1.8.1 ShuffleMapStage与ResultStage

     ShuffleMapStage的结束伴随着shuffle文件的写磁盘。

     ResultStage基本上对应代码中的action算子,即将一个函数应用在RDD的各个partition的数据集上,意味着一个job的运行结束

1.8.2 Shuffle中的任务个数

    Spark Shuffle分为map阶段和reduce阶段,或者称之为ShuffleRead阶段和ShuffleWrite阶段

1.8.2.1 map任务个数

    初始RDD分区个数由该文件的split个数决定,一个split对应生成的RDD的一个partition,初始RDD经过一系列算子计算后(假设没有执行repartition和coalesce算子进行重分区,则分区个数不变,仍为N,如果经过重分区算子,那么分区个数变为M),当执行到Shuffle操作时,map端的task个数和partition个数一致,即map task为N个。

1.8.2.2 reduce任务个数

    reduce端的stage默认取spark.default.parallelism这个配置项的值作为分区数,如果没有配置,则以map端的最后一个RDD的分区数作为其分区数(也就是N),分区数就决定了reduce端的task的个数

1.8.3 reduce端数据读取流程

    map端task和reduce端task不在相同的stage中,map task位于ShuffleMapStage,reduce task位于ResultStage

1.8.3.1 读取流程

     1)map task 执行完毕后会将计算状态以及磁盘小文件位置等信息封装到MapStatus对象中,然后由本进程中的MapOutPutTrackerWorker对象将mapStatus对象发送给Driver进程的MapOutPutTrackerMaster对象;

     2)在reduce task开始执行之前会先让本进程中的MapOutputTrackerWorker向Driver进程中的MapoutPutTrakcerMaster发动请求,请求磁盘小文件位置信息;

     3)当所有的Map task执行完毕后,Driver进程中的MapOutPutTrackerMaster就掌握了所有的磁盘小文件的位置信息。此时MapOutPutTrackerMaster会告诉MapOutPutTrackerWorker磁盘小文件的位置信息;

     4)完成之前的操作之后,由BlockTransforService去Executor0所在的节点拉数据,默认会启动五个子线程。每次拉取的数据量不能超过48M(reduce task每次最多拉取48M数据,将拉来的数据存储到Executor内存的执行(Execution)内存中)

1.8.4 shuffle类型

1.8.4.1 HashShuffle(1.6之前)

1)Shuffle write

    将 map 端划分数据、持久化数据的过程称为 shuffle write

(1) 未经优化的

    ShuffleMapTask在将数据写入磁盘之前,会先写入内存缓冲区,缓冲区被称为bucket,其大小为spark.shuffle.file.buffer.kb ,默认是 32KB(Spark 1.1 版本以前是 100KB)。每个 ShuffleMapTask 包含 R 个缓冲区,R = reducer 个数(也就是下一个 stage 中 task 的个数)。Task具体将数据写入哪一个bucket,由partitioner.partition(record.getKey()))决定,即对数据进行分区。当缓冲区溢出时,将数据刷写到磁盘上,一个bucket形成一个ShuffleBlockFile。由于每次shuffle使用M*R个bucket(M为MapTask个数,R为ReduceTask个数),会导致产生MR个shuffle文件,文件数量太多

(2)优化的Consolidate Shuffle

    Consolidate Shuffle 引入了shuffleFileGroup的概念,每个shuffleFileGroup都对应一批shuffle文件。shuffle文件数量与reduceTask数量相同(即与bucket数量相同)。只有在core上第一批执行的ShuffleMapTasks会创建一个shuffleFIleGroup,将数据写入到对应shuffle文件。在该 core 上后续执行的 ShuffleMapTasks 会复用shuffleFIleGroup和shuffle文件,即数据会继续写入到已有的shuffle文件。该机制会允许同个core上不同task复用同一个shuffle文件,对于多个task进行了一定程度的合并。这样,每次shuffle产生的文件数为C * R(C为spark集群的Core Number)。Consolidate Shuffle 功能可以通过spark.shuffle.consolidateFiles=true来开启。ShuffleMapTask在将数据写入bucket之前,根据mapSideCombine参数决定是否对数据进行combine操作,即map端的局部聚合。如果mapSideCombine为True,且指定了聚合函数,则会对数据先进行combine操作,再写入bucket。

2)Shuffle read

    将 reducer 读入数据、aggregate 数据的过程称为 shuffle read。

    shuffle read首先进行fetch操作,将shuffle文件fetch到本地机器上。fetch 来的 ShuffleFile 要先在内存做缓冲,Spark 规定这个缓冲界限不能超过 spark.reducer.maxMbInFlight,默认大小为 48MB。如果定义了聚集操作,通过聚集函数对fetch到缓冲区中的数据进行aggregator聚集操作,并且是边fetch边aggregator。Spark的聚集aggregator方式分为两种:只使用内存(必须保证有足够内存)和内存+磁盘。

  • 内存:spark.shuffle.spill = false就只用内存。内存使用的是AppendOnlyMap ,类似 Java 的HashMap。从缓冲中 deserialize 出来一个 <Key, Value> record,直接将其放进 HashMap 里面。如果该 HashMap 已经存在相应的 Key,那么直接进行 aggregate 操作,aggregate(hashMap.get(Key), Value),所以聚合操作必须是 commulative的。若在Map中没有查找到,则插入其中。
  • 内存+磁盘:使用的是ExternalAppendOnlyMap,其持有一个 AppendOnlyMap,开始过程与只使用内存模式相同,但是如果 AppendOnlyMap 快被装满时检查一下内存剩余空间是否可以够扩展,够就直接在内存中扩展。如果内存空间不足,ExternalAppendOnlyMap 将AppendOnlyMap进行排序(此排序按照key.hashcode进行排序,排序是为了之后的merge-aggregate过程,如果不排序,无法将多个文件merge)后 spill 到磁盘上生成spilledMap 文件,再重新 new 出来一个 AppendOnlyMap重复上述操作。当所有记录处理完毕之后,先对内存中的AppendOnlyMap进行排序,然后在对AppendOnlyMap和所有的spilledMap 文件进行全局 merge-aggregate。

1.8.4.2 SortShuffle(1.6及以后)

1)Shuffle write

(1)BypassMergeSortShuffleWriter

    BypassMergeSortShuffleWriter和Hash Shuffle中的HashShuffleWriter实现基本一致,唯一的区别在于,map端的多个输出文件会被汇总为一个文件,map端结果按照bucket顺序依次写入磁盘文件中,这么处理后,Shuffle生成的文件数显著减少了,同时还会生成indexFile文件,记录各个bucket在dataFile中的位置,用于后续reducer随机读取文件。

(2)触发机制

    shuffle map task数量小于spark.shuffle.sort.bypassMergeThreshold参数的值(默认200)。

    不是聚合类的shuffle算子。

(2)SortShuffleWriter

    数据会先写入一个内存数据结构中,此时根据不同的shuffle算子,可能选用不同的数据结构,并判断是否达到临界阈值,达到后先根据key进行排序,并分批spill到磁盘,最后将之前spill的临时磁盘文件进行merge,同事单独写一份索引文件,标识下游各个task的数据文件中的start offset和end offset。(写内存数据结构时,如果是reduceByKey这种聚合类的shuffle算子,那么会选用Map数据结构,一边通过Map进行聚合,一边写入内存;如果是join这种普通的shuffle算子,那么会选用Array数据结构,直接写入内存)。

(3)UnsafeShuffleWriter[了解] 

2)Shuffle read

    与1.8.4.1 HashShuffle的Shuffle read一致

1.9 内存管理

1.9.1 堆内堆外内存规划

    堆内内存受到JVM统一管理,堆外内存是直接向操作系统进行内存的申请和释放

1.9.1.1 堆内内存

    大小:由 Spark 应用程序启动时的 –executor-memory 或 spark.executor.memory 参数配置。

    Executor 内运行的并发任务共享 JVM 堆内内存,这些任务在缓存 RDD 数据和广播(Broadcast)数据时占用的内存被规划为存储(Storage)内存,而这些任务在执行 Shuffle 时占用的内存被规划为执行(Execution)内存,剩余的部分不做特殊规划,那些 Spark 内部的对象实例,或者用户定义的 Spark 应用程序中的对象实例,均占用剩余的空间

1.9.1.2 堆外内存

    1)通过配置 spark.memory.offHeap.enabled 参数启用,并由 spark.memory.offHeap.size 参数设定堆外空间的大小

    2)为了进一步优化内存的使用以及提高 Shuffle 时排序的效率,Spark 引入了堆外(Off-heap)内存,使之可以直接在工作节点的系统内存中开辟空间,存储经过序列化的二进制数据

    3)Executor的堆外内存主要用于程序的共享库、Perm Space、 线程Stack和一些Memory mapping等, 或者类C方式allocate object

1.9.2 内存空间分配

1.9.2.1 静态内存管理(1.6之前)

1)堆内内存

2)堆外内存

1.9.2.2 统一内存管理(1.6之后)

与静态内存管理的区别在于存储内存和执行内存共享同一块空间,可以动态占用对方的空闲区域

1)堆内内存

2)堆外内存

1.9.3 存储内存管理

1.9.3.1 RDD持久化机制

    堆内和堆外存储内存的设计,便可以对缓存 RDD 时使用的内存做统一的规划和管理,RDD 的每个 Partition 经过处理后唯一对应一个 Block,Driver端的Master 负责整个 Spark 应用程序的 Block 的元数据信息的管理和维护,而Executor端的 Slave 需要将 Block 的更新等状态上报到 Master,同时接收 Master 的命令。Spark 规定了 MEMORY_ONLY、MEMORY_AND_DISK 等 7 种不同的存储级别

1.9.3.2 RDD缓存

     1)RDD在缓存到存储内存之前,Partition中的数据一般以迭代器(Iterator)的数据结构来访问。通过迭代器可以获取分区中每一条序列化或者非序列化的数据项(Record),这些Record的对象实例在逻辑上占用了JVM堆内内存的other部分的空间,同一Partition的不同Record的空间并不连续。

     2)RDD在缓存到存储内存之后,Partition被转换成Block,Record在堆内或堆外存储内存中占用一块连续的空间。将Parititon由不连续的存储空间转换为连续存储空间的过程,Spark称之为展开(Unroll)。Block有序列化和非序列化两种存储格式,具体以哪种方式取决与该RDD的存储级别。每个Executor的Storage模块用一个链式Map结构(LinkedHashMap)来管理堆内和堆外存储内存中的所有Block对象的实例,对于这个LinkedHashMap新增和删除间接记录了内存的申请和释放。

     3)因为不能保证存储空间可以一次容纳Iterator中的所有数据,当前的计算任务在Unroll时要向MemeoryManager申请足够的Unroll空间来临时占位,空间不足则Unroll失败,空间足够时可以继续进行。对于序列化的Partition,其所需的Unroll空间可以直接累加计算,一次申请。而非序列化的Partition则要在遍历Record过程中依次申请,即每读取一条Record,采样估算其所需的Unroll空间进行申请,空间不足时可以中断,释放已占用的Unroll空间。如果最终Unroll成功,当前Partition所占用的Unroll空间被转换为正常缓存RDD的存储空间。

1.9.3.3 淘汰与落盘

    由于同一个Executor的所有的计算任务共享有限的存储内存空间,当有新的Block需要缓存但是剩余空间不足且无法动态占用时,就要对LinkedHashMap中的旧Block进行淘汰(Eviction),而被淘汰的Block如果其存储级别中同时包含存储到磁盘的要求,则要对其进行落盘(DROP),否则直接删除该Block

1.9.4 执行内存管理

1.9.4.1 Shuffle write

     1)若在map端选择普通的排序方式,会常用ExternalSorter进行排序,在内存中存储数据时主要占用堆内执行空间。

     2)若在map端选择Tungsten的排序方式,则采用ShuffleExternalSorter直接以序列化形式存储的数据排序,在内存中存储数据时可以占用堆外或堆内执行空间,取决于用户是否开启了堆外内存以及堆外执行内存是否足够。

1.9.4.2 Shuffle read

     1)在对reduce端的数据进行聚合时,要将数据交给Aggregator处理,在内存中存储数据时占用堆内执行空间。

     2)如果需要进行最终结果排序,则要再次将数据交给ExternalSorter处理,占用堆内执行空间。

     3)在ExternalSorter和Aggreator中,Spark会使用一种叫AppendOnlyMap的哈希表在堆内执行内存中存储数据,但在Shuffle过程中所有数据并不能都保存该Hash表中,当这个Hash表占用的内存会进行周期性采样,当其大到一定程度,无法再从MemoryManager申请到新的执行内存时,Spark就会将其全部内容存储到磁盘文件中,这个过程被称为溢存(Spill),溢存到磁盘的文件最后被归并(Merge)

1.10 BlockManger数据存储与管理机制

    BlockManager是整个Spark底层负责数据存储与管理的一个组件,Driver和Executor的所有数据都由对应的BlockManager进行管理。Driver上有BlockManagerMaster,负责对各个节点上的BlockManager内部管理的数据的元数据进行维护,比如block的增删改等操作,都会在这里维护好元数据的变更。 每个节点都有一个BlockManager,每个BlockManager创建之后,会向BlockManagerMaster进行注册,此时BlockManagerMaster会创建对应的BlockManagerInfo。

1.11 共享变量底层实现

1.11.1 广播变量

    允许编程者在每个Executor上保留外部数据的只读变量,而不是给每个任务发送一个副本,减少变量到各个节点的网络传输消耗,以及在各个节点上的内存消耗。在任务运行时,Executor并不获取广播变量,当task执行到 使用广播变量的代码时,会向Executor的内存中请求广播变量

1.11.2 累加器

    仅仅被相关操作累加的变量,Accumulator是存在于Driver端的,集群上运行的task进行Accumulator的累加,随后把值发到Driver端,在Driver端汇总,由于Accumulator存在于Driver端,从节点读取不到Accumulator的数值。

1.12 通讯架构

1.12.1 概述

    Spark2.x版本使用Netty通讯框架作为内部通讯组件。基于netty新的rpc框架借鉴了Akka的中的设计,基于Actor模型,spark各个组件可认为是独立的实体,实体间通过消息来通讯,EndPoint有1个InBox和N个OutBox(N>=1,取决于该EndPoint与多少个其他EndPoint通讯),接收的消息写入InBox,发送的消息写入OutBox

1.12.2 架构解析


    1) RpcEndpoint:RPC端点,Spark针对每个节点(Client/Master/Worker)都称之为一个Rpc端点,且都实现RpcEndpoint接口,内部根据不同端点的需求,设计不同的消息和不同的业务处理,如果需要发送(询问)则调用Dispatcher;

    2) RpcEnv:RPC上下文环境,每个RPC端点运行时依赖的上下文环境称为RpcEnv;

    3) Dispatcher:消息分发器,针对于RPC端点需要发送消息或者从远程RPC接收到的消息,分发至对应的指令收件箱/发件箱。如果指令接收方是自己则存入收件箱,如果指令接收方不是自己,则放入发件箱;

    4) Inbox:指令消息收件箱,一个本地RpcEndpoint对应一个收件箱,Dispatcher在每次向Inbox存入消息时,都将对应EndpointData加入内部ReceiverQueue中,另外Dispatcher创建时会启动一个单独线程进行轮询ReceiverQueue,进行收件箱消息消费;

    5) RpcEndpointRef:RpcEndpointRef是对远程RpcEndpoint的一个引用。当我们需要向一个具体的RpcEndpoint发送消息时,一般我们需要获取到该RpcEndpoint的引用,然后通过该应用发送消息。

    6) OutBox:指令消息发件箱,对于当前RpcEndpoint来说,一个目标RpcEndpoint对应一个发件箱,如果向多个目标RpcEndpoint发送信息,则有多个OutBox。当消息放入Outbox后,紧接着通过TransportClient将消息发送出去。消息放入发件箱以及发送过程是在同一个线程中进行;

    7) RpcAddress:表示远程的RpcEndpointRef的地址,Host + Port。

    8) TransportClient:Netty通信客户端,一个OutBox对应一个TransportClient,TransportClient不断轮询OutBox,根据OutBox消息的receiver信息,请求对应的远程TransportServer;

    9) TransportServer:Netty通信服务端,一个RpcEndpoint对应一个TransportServer,接受远程消息后调用Dispatcher分发消息至对应收发件箱。


 

文章转载自BigData Scholar,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论