微信公众号搜"智元新知"关注
微信扫一扫可直接关注哦!

spark任务提交流程源码分析

我这里使用spark2.4.4版本;

1.入口脚本与入口类

在这里插入图片描述

进入org.apache.spark.deploy.SparkSubmit类的main方法

override def main(args: Array[String]): Unit = {
    val submit = new SparkSubmit() {
      self =>

      override protected def parseArguments(args: Array[String]): SparkSubmitArguments = {
        new SparkSubmitArguments(args) {
          override protected def logInfo(msg: => String): Unit = self.logInfo(msg)

          override protected def logWarning(msg: => String): Unit = self.logWarning(msg)
        }
      }

      override protected def logInfo(msg: => String): Unit = printMessage(msg)

      override protected def logWarning(msg: => String): Unit = printMessage(s"Warning: $msg")

      override def doSubmit(args: Array[String]): Unit = {
        try {
          super.doSubmit(args)
        } catch {
          case e: SparkUserAppException =>
            exitFn(e.exitCode)
        }
      }

    }
	//提交
    submit.doSubmit(args)
  }

下一步

  def doSubmit(args: Array[String]): Unit = {
    // Initialize logging if it hasn't been done yet. Keep track of whether logging needs to
    // be reset before the application starts.
    val uninitLog = initializeLogIfNecessary(true, silent = true)

    val appArgs = parseArguments(args)
    if (appArgs.verbose) {
      logInfo(appArgs.toString)
    }
    appArgs.action match {
      //提交
      case SparkSubmitaction.SUBMIT => submit(appArgs, uninitLog)
      //杀掉
      case SparkSubmitaction.KILL => kill(appArgs)
      case SparkSubmitaction.REQUEST_STATUS => requestStatus(appArgs)
      case SparkSubmitaction.PRINT_VERSION => printVersion()
    }
  }

版权声明:本文内容由互联网用户自发贡献,该文观点与技术仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 [email protected] 举报,一经查实,本站将立刻删除。

相关推荐