【发布时间】:2017-12-25 14:36:25
【问题描述】:
例如我需要获取所有可用执行器的列表及其各自的多线程容量(不是总多线程容量,sc.defaultParallelism 已经处理了)。
由于这个参数是依赖于实现的(YARN 和 spark-standalone 有不同的分配核的策略)和情景(它可能因为动态分配和长期作业运行而波动)。我不能用其他方法来估计这个。有没有办法在分布式转换中使用 Spark API 检索这些信息? (例如TaskContext、SparkEnv)
更新至于Spark 1.6,我尝试了以下方法:
1) 运行具有多个分区 (>> defaultParallelism) 的 1 阶段作业,并计算每个 executorID 的不同线程 ID 的数量:
val n = sc.defaultParallelism * 16
sc.parallelize(n, n).map(v => SparkEnv.get.executorID -> Thread.currentThread().getID)
.groupByKey()
.mapValue(_.distinct)
.collect()
然而,这会导致估计值高于实际的多线程容量,因为每个 Spark 执行器都使用过度配置的线程池。
2) 类似于 1,除了 n = defaultParallesim,并且在每个任务中我添加一个延迟以防止资源协商器不平衡分片(快速节点完成它的任务并在慢速节点开始运行之前要求更多):
val n = sc.defaultParallelism
sc.parallelize(n, n).map{
v =>
Thread.sleep(5000)
SparkEnv.get.executorID -> Thread.currentThread().getID
}
.groupByKey()
.mapValue(_.distinct)
.collect()
它大部分时间都有效,但比必要的慢得多,并且可能会被非常不平衡的集群或任务推测破坏。
3)这个我没试过:使用java反射读取BlockManager.numUsableCores,这显然不是一个稳定的方案,内部实现随时可能发生变化。
如果你发现了更好的东西,请告诉我。
【问题讨论】:
-
谢谢 Paul,这是给 scala 的,我在深夜发布,所以没有写下我的调查,以后再补充
-
@Paul 更新了,够好吗?
-
看起来比原来好多了。
-
您可以查看
SparkContext.getLocalProperty,尤其是spark.executor.cores属性。见spark.apache.org/docs/latest/configuration.html -
@Adonis,这是核心数量的预期上限,而不是真实的。
标签: multithreading scala apache-spark distributed-computing