【问题标题】:Method to get number of cores for a executor on a task node?获取任务节点执行程序的核心数的方法?
【发布时间】: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


【解决方案1】:

使用 Spark REST API 非常简单。您必须获取应用程序 ID:

val applicationId = spark.sparkContext.applicationId

用户界面网址:

val baseUrl = spark.sparkContext.uiWebUrl

和查询:

val url = baseUrl.map { url => 
  s"${url}/api/v1/applications/${applicationId}/executors"
}

使用 Apache HTTP 库(已在 Spark 依赖项中,改编自 https://alvinalexander.com/scala/scala-rest-client-apache-httpclient-restful-clients):

import org.apache.http.impl.client.DefaultHttpClient
import org.apache.http.client.methods.HttpGet
import scala.util.Try

val client = new DefaultHttpClient()

val response = url
  .flatMap(url => Try{client.execute(new HttpGet(url))}.toOption)
  .flatMap(response => Try{
    val s = response.getEntity().getContent()
    val json = scala.io.Source.fromInputStream(s).getLines.mkString
    s.close
    json
  }.toOption)

和json4s:

import org.json4s._
import org.json4s.jackson.JsonMethods._
implicit val formats = DefaultFormats

case class ExecutorInfo(hostPort: String, totalCores: Int)

val executors: Option[List[ExecutorInfo]] = response.flatMap(json => Try {
  parse(json).extract[List[ExecutorInfo]]
}.toOption)

只要您保留应用程序 ID 和 ui URL 并打开 ui 端口到外部连接,您就可以在任何任务中执行相同的操作。

【讨论】:

  • 非常感谢您的回答!再等几个星期吧,我觉得如果不小心使用,这可能会成为一种反模式,一个 spark master 可能会管理数千个节点,它的 UI 并不是设计为通过一个效率不高的数据被所有人 DDoSed序列化协议。
【解决方案2】:

我会尝试以类似于 Web UI 的方式实现 SparkListener。 This code 作为示例可能会有所帮助。

【讨论】:

  • 好主意!在 Spark 1.6 中,这是 ExecutorInfo 唯一可读的地方,所以也许值得一试。唯一的缺点是监听器只在驱动上触发,所以它的执行不是本地的。
猜你喜欢
  • 1970-01-01
  • 2013-05-27
  • 2018-07-09
  • 2019-11-01
  • 1970-01-01
  • 1970-01-01
  • 2014-08-28
  • 1970-01-01
  • 2019-10-18
相关资源
最近更新 更多