Scala Spark独立Stage未按预期并行执行问题排查
Spark独立Stage串行执行问题排查方案
问题现象
排查Spark作业运行问题时,从Spark UI的DAG视图可见应用共包含4个Stage,DAG结构如下:
0 1 2 \ | / \\ | / 3
- 预期行为:前3个无依赖的Stage(Stage0、Stage1、Stage2)并行执行,三者全部完成后再启动依赖三者输出的Stage3
- 实际观测:三个Stage同时被提交,Stage3处于pending状态时三个Stage在UI均显示为active,但并未真正并行运行——Stage1需等待Stage0执行完成才启动,Stage2需等待Stage1执行完成才启动。
关联核心代码
def runJob(ss: SparkSession): Unit = { (records _) .andThen(convertToOtherFormat) .andThen(writeRecords(_, path)) } def records: Dataset[Record] = { import sparkSession.implicits._ region.cityIds // Set[cityId] .map(getContentDataframe) // Set[sql.DataFrame] .reduce(_ union _) // sql.DataFrame .coalesce(numPartitions) // Dataset[Row] .as[Record] // Dataset[Record] } def convertToOtherFormat(records: Dataset[_ <: Record]): RDD[(Key, Value)] = { records // DAG中该行标记为Stage0/1/2的起始位置 .rdd .map(record => { add(record) ( new Key(record.key), new Value(record.value) ) }) }
常见认知误区澄清
首先纠正几个排查过程中的错误判断:
- Spark UI的active状态不代表正在运行:Stage标记为active仅代表已经被提交到调度池等待资源分配,是否实际运行要看当前是否有空闲计算资源、调度器是否给它分配资源。
- Scala集合的
.par并行不等于Spark作业并行:.par实现的是驱动端JVM内的多线程集合遍历,你在构造DataFrame时加.par,只是并行生成DataFrame对象——DataFrame本身只是逻辑执行计划的封装,构造过程不会触发任何分布式计算,因此完全不会改变Stage的调度顺序。 - 单作业内的union分支不会自动并行计算:你写的
map(getContentDataframe).reduce(_ union _)只是把多个分支的执行计划拼接成一个大的执行计划,最终触发action时属于同一个Spark作业,默认不会自动并行跑多个分支。
排查&修复步骤
按优先级依次排查以下问题:
1. 先核对计算资源是否足够
先到Spark UI的Executors页面计算总可用计算core:总core数 = executor实例数 * 单executor配置的core数,再对比单个Stage的总task数(和对应数据集的分片数一致)。
- 如果总core数 ≤ 单个Stage的task并发数,就算调度器想同时跑多个Stage,也会因为没有空闲core导致后续Stage排队,表现出来就是串行执行。这种情况先调大executor资源配置,保证总core数至少能同时容纳2~3个Stage的task并发运行。
2. 调整调度器模式
如果资源足够但还是串行,检查Spark调度模式配置,默认的FIFO(先进先出)调度器会优先给先提交的Stage分配全部资源,后续Stage必须等前面的Stage释放资源才能运行。将调度模式改为FAIR(公平调度)即可让多个active Stage公平共享资源:
spark.conf.set("spark.scheduler.mode", "FAIR")
3. 强制多分支并行(如果上述两步配置后仍不满足需求)
如果需要严格保证三个分支完全并行,不要直接在串行构造的union结果上触发action,而是在驱动端用独立线程池分别提交各个分支的计算,物化缓存后再做union:
import scala.concurrent.{ExecutionContext, Future, Await} import scala.concurrent.duration._ // 线程池大小匹配需要并行的分支数量 implicit val ec: ExecutionContext = ExecutionContext.fromExecutorService( java.util.concurrent.Executors.newFixedThreadPool(3) ) val cityDataframes = region.cityIds.map(getContentDataframe) // 并行提交各分支计算,通过cache物化结果 val futureCachedDFs = Future.sequence( cityDataframes.map(df => Future { val cachedDF = df.cache() cachedDF.count() // 触发action完成实际计算 cachedDF }) ) val cachedDataframes = Await.result(futureCachedDFs, 1.hour) // 所有分支计算完成后再做union和后续逻辑 val finalRecords = cachedDataframes.reduce(_ union _) .coalesce(numPartitions) .as[Record]
使用该方式时必须保证FAIR调度已开启、集群资源足够支撑三个作业同时运行。
学习参考方向
- Spark核心调度原理:重点掌握作业/Stage/Task的提交流程、FIFO与FAIR调度器的实现差异
- Spark Dataset执行计划逻辑:理解转换算子(如
union)与行动算子的区别,掌握执行计划生成、物化的时机 - 驱动端并行提交作业逻辑:明确区分驱动端多线程并行与Executor端分布式计算并行的差异,不要混淆Scala本地并行与Spark分布式并行的适用场景
内容的提问来源于stack exchange,提问作者lslslsls
相关产品推荐
相关产品推荐

