ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

深入解析Spark Stage划分与Pipeline计算:DAG调度核心原理与实践优化

深入解析Spark Stage划分与Pipeline计算:DAG调度核心原理与实践优化 1. 从“一个任务”到“多个阶段”Spark作业执行的核心逻辑当我们把一个复杂的计算任务提交给Spark时它并不会一股脑地开始执行。想象一下你要做一顿大餐你不会把所有食材同时扔进锅里。你会先洗菜、切菜、备料这些步骤可以并行然后开火炒菜这依赖于备好的料最后装盘上桌。Spark处理数据任务的方式与此高度相似它将一个庞大的作业Job拆分成多个可并行、有依赖关系的计算阶段Stage这种精密的拆分机制是Spark实现高效、容错分布式计算的核心。这个拆分过程的核心“导演”是DAGScheduler。它站在全局视角审视由用户代码一系列RDD转换操作构成的有向无环图DAG。DAGScheduler的工作就是分析这个DAG找出哪些转换可以“流水线”式地合并执行这就是Pipeline计算模式而哪些转换必须让所有并行任务停下来进行数据的重新洗牌和交换这就是Shuffle也是Stage划分的依据。最终它将DAG切分成一个个Stage交给TaskScheduler去调度成具体的Task在集群上运行。理解Stage的划分与计算模式远不止是记忆概念。它能直接帮助你性能调优一眼看出作业的Shuffle次数这是性能瓶颈的主要来源。故障排查当任务失败时能快速定位是哪个Stage出了问题该Stage的输入输出是什么。资源预估理解每个Stage的并行度从而合理设置资源。代码优化主动设计RDD转换链减少不必要的Shuffle构建更高效的Pipeline。接下来我们就深入这个“后厨”看看DAGScheduler是如何阅读“菜谱”DAG并制定出高效的“烹饪阶段”Stage计划的。2. DAGScheduler的视角如何解读RDD的血缘与依赖要理解Stage的划分首先要站在DAGScheduler的角度看看它眼里的用户代码是什么样子。我们写的map、filter、groupByKey等操作在Spark内部被翻译成一系列RDD及它们之间的依赖关系。这种依赖关系被称为“血统”Lineage它不仅是容错恢复的日志更是Stage划分的蓝图。DAGScheduler主要关注两种依赖关系这是划分Stage的基石2.1 窄依赖Narrow Dependency构建Pipeline的基石窄依赖指的是父RDD的每个分区最多被一个子RDD的分区所依赖。这是一种“一对一”或“一对固定个”的映射关系。典型操作map、filter、flatMap、union等。特点数据局部性计算子RDD分区所需的所有数据都来自于父RDD的某一个或某几个确定的分区。这意味着数据不需要在网络间移动。管道化Pipelining执行正因为数据不需要移动多个窄依赖转换操作可以被合并到一个Stage内形成一个连续的计算管道。在一个Task内数据可以像在流水线上一样依次经过map、filter等操作而无需将中间结果落盘或进行网络传输效率极高。高效容错如果某个分区数据丢失只需要根据血统重新计算其父RDD的对应分区即可无需回溯整个RDD。示例rdd.map(func1).filter(func2)map和filter之间是窄依赖。DAGScheduler会将这两个操作合并到同一个Stage中。一个Task会连续地对一个数据分片先执行func1再执行func2。2.2 宽依赖Shuffle DependencyStage的边界宽依赖指的是父RDD的一个分区可能被多个子RDD的分区所依赖。这是一种“一对多”的关系它要求所有父RDD分区的数据按照某种规则重新洗牌和分发。典型操作groupByKey、reduceByKey、join非相同分区器时、repartition、coalesce增加分区数时等。特点需要Shuffle这是宽依赖最核心的标志。它要求上游Stage的所有Task必须先将自己的计算结果按照Key进行分区、排序可选并持久化到本地磁盘然后下游Stage的Task再从各个节点拉取Fetch属于自己的那部分数据。这个过程涉及大量的磁盘I/O和网络I/O是Spark作业中最耗时、最昂贵的操作。划分Stage宽依赖是Stage划分的天然边界。DAGScheduler在遇到宽依赖时就会在此处“切一刀”将宽依赖之前的操作划分为一个Stage称为ShuffleMapStage其Task产出Shuffle数据宽依赖之后的操作划分为下一个Stage称为ResultStage其Task产出最终结果或下一个Shuffle的输入。容错成本高如果下游Stage的数据丢失由于它依赖上游所有分区的数据通常需要重新计算整个上游Stage。示例rdd.map(...).reduceByKey(...).map(...)reduceByKey是一个宽依赖。DAGScheduler会在这里进行切分Stage 0: 包含第一个map操作是一个ShuffleMapStage。它的Task对数据进行处理并为reduceByKey准备Shuffle输出。Stage 1: 包含reduceByKey和最后一个map操作。这里reduceByKey本身需要Shuffle读而它和最后的map之间又是窄依赖所以被合并到同一个ResultStage中。注意coalesce操作在减少分区数时例如从100个分区合并到10个可能不需要全量Shuffle因为多个父分区可以合并到一个子分区此时是窄依赖。但在增加分区数时必然需要Shuffle来创建新分区此时是宽依赖。这是一个需要根据参数仔细辨别的特例。3. Stage划分算法DAGScheduler的“切图”策略了解了窄依赖和宽依赖我们就可以梳理出DAGScheduler划分Stage的具体算法了。这个过程本质上是在RDD DAG中进行反向解析遇到宽依赖就建立一个新的Stage。算法步骤详解从最终RDD出发DAGScheduler从触发行动操作如collect,saveAsTextFile,count所产生的最终RDD开始回溯。递归回溯遇窄则并沿着RDD的血缘关系依赖链反向递归。如果当前依赖是窄依赖则将这个RDD和它所依赖的父RDD划归到同一个Stage中。这个过程会一直进行尽可能多地将连续的窄依赖操作合并。遇宽则切注册阶段当回溯遇到一个宽依赖时DAGScheduler会执行以下操作切断在此宽依赖处“切一刀”。宽依赖之前的RDD操作链属于当前正在构建的Stage我们称之为stageA。注册将stageA注册为一个独立的Stage通常是一个ShuffleMapStage除非它是最后一个Stage。回溯然后将这个宽依赖所依赖的父RDD作为新的起始点重复步骤2和3开始构建下一个StagestageB。循环直至源头重复步骤2和3直到回溯到最源头的RDD如从HDFS读取数据生成的RDD。这个最源头的RDD会被划入第一个Stage。确定Stage类型如果一个Stage的输出需要为后续的Shuffle提供数据它就是ShuffleMapStage。它的每个TaskShuffleMapTask会生成一个数据文件供下游Stage拉取。如果一个Stage直接生成最终结果即其后没有其他Stage它就是ResultStage。它的每个TaskResultTask直接计算并返回结果给Driver或写入外部存储。用一个复杂例子串联理解 假设我们有如下代码val rddA sc.textFile(...) val rddB rddA.map(...).filter(...) val rddC rddB.groupByKey(...) val rddD rddC.map(...) val rddE rddD.join(rddB) val result rddE.reduceByKey(...).collect()划分过程从最终result对应rddE.reduceByKey(...)回溯。遇到reduceByKey宽依赖切分。reduceByKey及其之前的rddE属于Stage 3ResultStage。回溯rddE它由rddD.join(rddB)产生。join通常也是宽依赖除非两者有相同分区器且分区数相同。在此处切分。join操作本身及其左分支rddD、右分支rddB的依赖链需要被包含进Stage 2ShuffleMapStage。注意rddB的链rddA.map.filter是窄依赖被合并进来。回溯rddD它由rddC.map(...)产生窄依赖继续回溯。回溯rddC它由rddB.groupByKey(...)产生。groupByKey是宽依赖切分。groupByKey及其之前的rddB链rddA.map.filter属于Stage 1ShuffleMapStage。回溯rddB的源头rddA它是从文件读取生成的回溯结束。最终划分Stage 0: 读取文件生成rddA。如果文件读取后直接是rddA且没有其他依赖它可能被合并到Stage 1中但逻辑上它是起点Stage 1:rddA.map.filter.groupByKey的前半部分到Shuffle写为止是一个ShuffleMapStage。Stage 2:rddC.map.join是一个ShuffleMapStage因为它的输出供Stage 3的reduceByKey做Shuffle读。Stage 3:rddE.reduceByKey是最终的ResultStage。通过这个划分DAGScheduler将复杂的DAG转化为了一个Stage的有向无环图每个Stage内部都是高效的Pipeline计算Stage之间则通过Shuffle进行数据交换。4. Pipeline计算模式Stage内部的“流水线”加速前面提到Stage内部是连续的窄依赖变换它们会以Pipeline模式执行。这是Spark提升CPU和内存利用率的关键优化。Pipeline的本质将多个函数如f,g,h组合成一个函数h(g(f(x)))在一个Task内对数据分片的一次遍历中连续执行而不是每个函数都产生一个中间RDD物化在内存或磁盘上。工作原理函数融合Fusion在生成Stage的Task时Spark会将这个Stage内所有窄依赖操作对应的函数串联成一个单一的复合函数。迭代器模式Task从数据源或上一个Shuffle结果读取一个分区的数据这个数据通常以迭代器的形式提供。链式调用数据记录一条一条地通过这个复合函数迭代器。对于一条记录先执行第一个map的f函数其结果立即作为输入传给filter的g函数再传给下一个map的h函数。处理完一条记录再处理下一条。全内存操作整个过程中间结果除了函数本身的临时对象不会以集合形式驻留在内存中更不会落盘极大地减少了GC压力和I/O开销。示例与对比 假设rdd.map(f).filter(g).map(h)无Pipeline低效模式Task先对整个分片数据执行map(f)生成一个完整的中间结果集A写入内存。然后对A执行filter(g)生成完整的中间结果集B写入内存。最后对B执行map(h)得到最终结果C。缺点产生了两个完整的中间数据集A和B占用双倍内存增加了GC频率。有Pipeline高效模式Task创建一个复合函数record h(g(f(record)))。从数据源读取一条记录record1计算h(g(f(record1)))输出结果result1。立即读取下一条记录record2计算h(g(f(record2)))输出result2。如此循环直到处理完所有记录。优点内存中几乎只同时存在一条记录的处理状态内存效率极高CPU缓存友好。Pipeline的边界Pipeline的边界就是Stage的边界。一旦遇到需要Shuffle的宽依赖当前Pipeline就必须中断因为所有分区数据必须完成计算、持久化并重新组织后才能进行下一步。因此优化Spark作业的一个核心原则就是通过调整算子或使用特定API如reduceByKey替代groupByKey尽可能减少Shuffle次数从而构建更长的、更高效的Pipeline。5. Shuffle详解Stage间数据重组的代价与优化Shuffle是连接不同Stage的桥梁也是分布式计算中最昂贵的操作。理解Shuffle的细节对于诊断性能问题和进行调优至关重要。Shuffle的两个阶段 以reduceByKey为例它触发一个Shuffle。这个过程涉及两个StageShuffle Write上游Stage的Task每个ShuffleMapTask将其计算结果Key, Value对在内存中进行聚合如果指定了mapSideCombine然后根据目标分区器Partitioner默认为HashPartitioner计算每个Key应该属于下游的哪个分区。接着它会为下游的每一个分区在本地磁盘上创建一个临时文件并将属于该分区的数据写入对应的文件。最后它会生成一个索引文件记录每个输出分区数据在文件中的起始和结束位置。这个索引文件会报告给Driver的DAGScheduler。Shuffle Read下游Stage的Task每个ResultTask或下一个ShuffleMapTask启动时会向Driver询问自己需要的数据位于哪些上游节点。然后它向这些节点发起网络请求拉取Fetch属于自己的那部分数据文件。拉取到的数据可能在内存中进行合并、排序如果指定了和聚合如reduceByKey的reduce函数最终形成这个Task的输入数据。Shuffle的潜在问题与调优方向数据倾斜某个Key对应的数据量远大于其他Key导致某个下游Task处理的数据量巨大成为最慢的“拖尾任务”。应对使用sample算子探查数据分布考虑使用加盐Salting技术将大Key打散对于聚合操作尝试使用reduceByKey而非groupByKey因为前者能在Map端先进行局部合并大幅减少Shuffle数据量。Shuffle文件数量爆炸如果上游Stage有M个分区下游Stage有N个分区那么最坏情况下会产生 M * N 个Shuffle临时文件。虽然Spark后续进行了优化如将同一分区数据合并到一个文件但分区数过多仍会带来大量随机I/O和小文件问题。应对合理设置分区数。通过spark.sql.shuffle.partitionsSQL或rdd.repartition/coalesceRDD来控制。分区数并非越多越好一般建议为集群总核心数的2-3倍。磁盘I/O与网络I/OShuffle数据需要落盘防止Executor内存溢出并在网络间传输。应对使用高性能磁盘如SSD作为spark.local.dir在万兆网络环境下运行启用Shuffle压缩spark.shuffle.compress以减少网络传输量。内存压力Shuffle Read阶段拉取的数据需要在内存中进行缓冲、合并和聚合。应对调整spark.shuffle.memoryFractionSpark 1.x或spark.memory.fraction中用于Shuffle的比例对于非常复杂的聚合考虑增加Executor内存或调整聚合算法的内存使用方式。一个关键配置mapSideCombine在reduceByKey等聚合操作中mapSideCombine默认为true是一个极其重要的优化。它意味着在Shuffle Write之前先在每个Mapper上游Task内部对本地相同的Key进行一次预聚合Combine。这能显著减少需要Shuffle的数据量。相比之下groupByKey不会进行Map端聚合它会原样传输所有数据因此通常比reduceByKey更慢、更耗资源。在代码中应优先使用reduceByKey、aggregateByKey等可进行Map端合并的算子。6. 实战洞察从UI与日志中观察Stage与Shuffle理论需要结合实践。Spark Web UI和日志是我们观察Stage划分与Shuffle行为的最佳窗口。在Spark Web UI中Jobs Stages标签页提交一个作业后这里会清晰显示生成了几个Job每个Job又被划分成几个Stage。Stage以DAG图或列表形式展示依赖关系一目了然。每个Stage框会显示其ID、描述、提交时间、持续时间、Task数量等信息。识别Shuffle在Stage图中Stage之间的箭头连线就代表Shuffle依赖。鼠标悬停可以看到Shuffle的数据量Shuffle Read/Write Size。数据量异常大的Shuffle箭头就是性能瓶颈的明显信号。Stage详情点击一个Stage可以进入其详情页。这里可以看到所有Task的执行情况、时间分布、GC时间、Shuffle数据量等。如果某个Stage的某些Task执行时间远长于其他Task很可能就是遇到了数据倾斜。Event Timeline可以可视化地看到各个Executor上Task的执行时间线Shuffle Read/Write的时间段也会被标记出来帮助你直观了解资源利用情况和Shuffle耗时。在日志中 在Driver的日志或stdout中DAGScheduler会打印关键信息INFO DAGScheduler: Submitting ShuffleMapStage 0 (...), which has no missing parents INFO DAGScheduler: Submitting ResultStage 1 (...) INFO DAGScheduler: ShuffleMapStage 0 finished in 2.3 s INFO DAGScheduler: looking for newly runnable stages INFO DAGScheduler: running: Set(ShuffleMapStage 0) INFO DAGScheduler: waiting: Set(ResultStage 1) INFO DAGScheduler: failed: Set()这些日志清晰地展示了Stage的提交、执行和依赖关系。当作业失败时日志会明确指出是哪个Stage的哪个Task失败了以及失败原因这是排查故障的第一现场。调试技巧当你对Stage划分有疑问时可以在代码中对RDD调用toDebugString方法它会打印出RDD的血缘关系图其中用缩进和-清晰地显示了窄依赖和宽依赖ShuffleDependency。使用rdd.dependencies属性可以查看RDD的直接依赖列表判断是NarrowDependency还是ShuffleDependency。理解Stage的划分与计算模式是掌握Spark性能调优的“内功心法”。它让你从“黑盒使用”变为“白盒掌控”能够预见代码的执行路径精准定位瓶颈所在。记住核心链条代码 - RDD DAG - (宽依赖切分) - Stages - (Pipeline执行) - Tasks - Shuffle (数据交换)。在编写下一个Spark作业时不妨先在脑中推演一下它的Stage图思考如何减少图中的“宽依赖箭头”这将是你写出高效Spark代码的关键一步。
返回列表