【问题标题】:How can I get Apache Spark job progress without knowing the job id?如何在不知道作业 ID 的情况下获得 Apache Spark 作业进度?
【发布时间】:2016-01-29 23:25:06
【问题描述】:

我举个例子:

val sc: SparkContext // An existing SparkContext.
val sqlContext = new org.apache.spark.sql.SQLContext(sc)

val df = sqlContext.read.json("examples/src/main/resources/people.json")

df.count

我知道我可以使用 Spark 上下文使用 SparkListener 监控作业;但是,这会给我所有工作的事件(我不能使用,因为我不知道工作 I​​D)。

如何才能仅获取“计数”操作的进度?

【问题讨论】:

  • 您需要以编程方式对其进行监控吗?
  • 是的。因为我会将这些指标传递给其他应用程序。
  • 你考虑过解析网页界面吗?在描述下,您可以查找df.count,行条目中的其他字段应为您提供持续时间、阶段:成功/总计和任务:成功/总计
  • 我认为 Web UI 是要走的路。它很容易通过描述性链接分为作业和阶段。见Spark UI

标签: apache-spark


【解决方案1】:

正如 cmets 中已经建议的那样,可以使用REST API of the Spark UI 来收集所需的数字。

主要问题是确定您感兴趣的阶段。从代码到阶段没有 1:1 的映射。例如,单个计数将触发两个阶段(一个阶段用于计算数据帧每个分区中的元素,第二个阶段用于汇总第一阶段的结果)。阶段通常获取触发其执行的操作的名称,尽管这可能是 changed within the code

可以创建一个方法来查询具有特定名称的所有阶段的 REST API,然后将这些阶段的所有任务数量以及已完成任务的数量相加。假设所有任务将大致花费相似的执行时间(如果数据集具有倾斜分区,则此假设是错误的),可以使用已完成任务的份额作为作业进度的衡量标准。

def countTasks(sparkUiUrl: String, stageName: String): (Int, Int) = {
  import scala.util.parsing.json._
  import scala.collection.mutable.ListBuffer
  def get(url: String) = scala.io.Source.fromURL(url).mkString

  //get the ids of all running applications and collect them in a ListBuffer
  val applications = JSON.parseFull(get(sparkUiUrl + "/api/v1/applications?staus=running"))
  val apps: ListBuffer[String] = new scala.collection.mutable.ListBuffer[String]
  applications match {
    case Some(l: List[Map[String, String]]) => l.foreach(apps += _ ("id"))
    case other => println("Unknown data structure while reading applications: " + other)
  }

  var countTasks: Int = 0;
  var countCompletedTasks: Int = 0;

  //get the stages for each application and sum up the number of tasks for each stage with the requested name
  apps.foreach(app => {
    val stages = JSON.parseFull(get(sparkUiUrl + "/api/v1/applications/" + app + "/stages"))
    stages match {
      case Some(l: List[Map[String, Any]]) => l.foreach(m => {
        if (m("name") == stageName) {
          countTasks += m("numTasks").asInstanceOf[Double].toInt
          countCompletedTasks += m("numCompleteTasks").asInstanceOf[Double].toInt
        }
      })
      case other => println("Unknown data structure while reading stages: " + other)
    }
  })

  //println(countCompletedTasks + " of " + countTasks + " tasks completed")
  (countTasks, countCompletedTasks)
}

为给定的计数示例调用此函数

println(countTasks("http://localhost:4040", "show at CountExample.scala:16"))

会打印出两个数字:第一个是所有任务的数量,第二个是完成任务的数量。

我已经用 Spark 2.3.0 测试了这段代码。在生产环境中使用它之前,它肯定需要一些额外的打磨,尤其是一些更复杂的错误检查。不仅可以通过统计已完成的任务,还可以通过统计失败的任务来改进统计数据。

【讨论】:

    猜你喜欢
    • 2019-11-01
    • 1970-01-01
    • 1970-01-01
    • 2021-07-29
    • 1970-01-01
    • 1970-01-01
    • 2017-07-02
    • 2017-05-12
    • 2021-08-12
    相关资源
    最近更新 更多