2026/7/28 11:36:11

深入理解Spark

深入理解Spark 文章目录1、Spark 是什么?2、Spark 运行模式3、适合Spark的场景4、Spark相关术语5、Spark程序执行流程6、理解Spark Stage的划分6.1 Spark Stage的划分6.2 Spark DAG的可视化7、Spark调度过程7.1 Spark的两级调度模型7.2 以Spark On Yarn说明调度过程小结在前面博客文章里,已经把大数据实时分析项目在spark组件之前的各个组件原理、部署和测试都给出相关讨论,接下来是项目最核心的内容:实时计算部分,因为项目将使用spark streaming做微批计算(准实时计算),因此接下的文章内容将深入spark以及spark streaming架构原理,为后面实际计算编程做铺垫。1、Spark 是什么?Spark是一种分布式的并行计算框架,什么是计算框架?所谓的计算(在数据层面理解)其实是用于数据处理和分析的一套解决方案,例如Python的Pandas,相信用过Pandas都很容易理解Pandas擅长做什么,加载数据、对数据进行各类加工、分析数据等,只不过Pandas只适合在单机上的、数据量百万到千万级的计算组件,而Spark则是分布式的、超大型多节点可并行处理数据的计算组件。Spark通常会跟MapReduce做对比,它与MapReduce 的最大不同之处在于Spark是基于内存的迭代式计算——Spark的Job处理的中间(排序和shuffling)输出结果可以保存在内存中,而不是在HDFS磁盘上反复IO浪费时间。除此之外,一个MapReduce 在计算过程中只有Map 和Reduce 两个阶段。而在Spark的计算模型中,它会根据rdd依赖关系预选设计出DAG计算图,把job分为n个计算阶段(Stage),因为它内存迭代式的,在处理完一个阶段以后,可以继续往下处理很多个阶段,而不只是两个阶段。Spark提供了超过80种不同的Transformation和Action算子,如map,reduce,filter,reduceByKey,groupByKey,sortByKey,foreach等,并且采用函数式编程风格,实现相同的功能需要的代码量极大缩小(尤其用Scala和Python写计算业务代码方面)。正是基于使用易用性,因此Spark能更好地用于基于分布式大数据的数据挖掘与机器学习等需要迭代的MapReduce的算法。Spark生态如下:2、Spark 运行模式目前最为常用的Spark运行模式有:Local:本地进程运行,例如启动一个pyspark交互式shell,一般用于开发调试Spark应用程序Standalone:利用Spark自带的资源管理与调度器运行Spark集群,采用Master/Slave结构,可引入ZooKeeper实现spark集群HAHadoop YARN : 集群运行在YARN资源管理器上,资源管理交给YARN,Spark只负责进行任务调度和计算,参考本博客《基于YARN HA集群的Spark HA集群》Apache Mesos :Apache Mesos abstracts CPU, memory, storage, and other compute resources away from machines (physical or virtual), enabling fault-tolerant and elastic distributed systems to easily be built and run effectively.将计算、内存、存储以及其他计算资源抽象出来,使得那些分布式系统更容易被构建以及更高效的被执行。其实就是跟YARN一样,做资源管理。Mesos和YARN两种资源有什么区别:之前看一个视频,对其给出的解释印象深刻:Mesos:细腻度资源管控YARN:粗粒度资源管控例如有个老师要给45个学生上课,向教务处申请课室资源,若教务处以Mesos模式发放资源,那么它会发放只能容纳45个学生的课室,典型的按需分配;若教务处以YARN模式发放资源,那么它会发放能容200个学生的大教室,但实际上还有155个人位置资源空闲。这就是资源的细腻度和粗粒度的区别。3、适合Spark的场景Spark适用场景:Spark是基于内存的迭代计算框架,适用于需要多次操作特定数据集的应用场合,基于大数据的机器学习再适合不过,例如梯度下降法,需要不断迭代找到全局或局部最优解。需要反复操作的次数越多,所需读取的数据量越大,受益越大,数据量小但是计算密集度较大的场合,受益就相对较小。准实时计算场合:实时接收用户行为原始数据,并通过Spark Streaming计算(转换+加工),例如在广告、报表、推荐系统等业务上,在广告业务方面需要大数据做应用分析、效果分析、定向优化等,在推荐系统方面则需要大数据优化相关排名、个性化推荐以及热点点击分析等,这些业务天生适合大型的互联网巨头。Spark不适用场景:内存消耗极大,在内存不足的情况下,Spark会下放到磁盘,会降低应有的性能有高实时性要求的流式计算业务,例如实时性要求毫秒级,对每一条数据都触发实时计算的,这种场合已经被Flink称霸。流线长或文件流量非常大的数据集不适合,这是因为这种场合rdd消耗极大的内存集群压力大时,一旦一个task失败会导致它前面一条线所有的前置任务全部重跑(尤其对于RDD 血缘关系链长且有多个宽依赖的情况),JVM的GC不够及时,内存不能及时释放,将会出现恶性循环导致更多的task失败,导致整个Application效率极低。所以为什么说Spark是适合“微批”处理,意味着每隔一段时间(1秒或者几秒不等)来一批次数据,Spark适合“一小口一小口准实时地吃数据”。4、Spark相关术语