我正在Spark 2.2上通过Spark submit在Mesos集群上运行我的工作。我宣布:
-
(在sparkSession配置中)spark。遗嘱执行人。核心5
-
(在spark submit配置中)执行器核心总数20
-
(在spark提交配置中)驱动核心20
所以我在一个驱动程序上有20个内核和4个执行器,每个执行器有5个内核。
通常当我看阶段进度条的时候
[Stage 12:================================> (86 + 20) / 100]
它表示已完成的任务数(86)、正在运行的任务数(20)和任务总数(100)。而且,在大多数情况下,正在运行的任务数不会超过中给出的值
total-executor-cores
所以在这个例子中是20。
然而现在我有一份工作,我看到了这样的东西
[Stage 22:================================> (355 + 155) / 568]
对于
相同的设置
如上所述,相同的集群。在SparkUI中,我看到每个执行器有5个内核(按要求),每个执行器一次运行大约30个任务,通常最多运行5个任务(每个内核1个)。我不明白为什么。
这可能与我的工作性质有关吗?它涉及将Spark数据帧转换为RDD,并在驱动程序上进行收集。示例阶段
val dfOnDriver = df.select("colA", "colB", "colC")
.rdd
.map(...)
.collect
-
问题1
为什么要运行这么多任务?
-
问题2
在基于执行者的正常计算中,有没有办法故意强迫这种行为?
-
问题3
这种行为是否需要,并且比每个核心1个任务更有效?
如果您需要有关我的设置的更多信息,请告诉我。