You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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)
      )
    })
}

常见认知误区澄清

首先纠正几个排查过程中的错误判断:

  1. Spark UI的active状态不代表正在运行:Stage标记为active仅代表已经被提交到调度池等待资源分配,是否实际运行要看当前是否有空闲计算资源、调度器是否给它分配资源。
  2. Scala集合的.par并行不等于Spark作业并行:.par实现的是驱动端JVM内的多线程集合遍历,你在构造DataFrame时加.par,只是并行生成DataFrame对象——DataFrame本身只是逻辑执行计划的封装,构造过程不会触发任何分布式计算,因此完全不会改变Stage的调度顺序。
  3. 单作业内的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.29 02:03:27