• linkedu视频
  • 平面设计
  • 电脑入门
  • 操作系统
  • 办公应用
  • 电脑硬件
  • 动画设计
  • 3D设计
  • 网页设计
  • CAD设计
  • 影音处理
  • 数据库
  • 程序设计
  • 认证考试
  • 信息管理
  • 信息安全
菜单
linkedu.com
导航菜单
  • 网页制作
  • 数据库
  • 程序设计
  • 操作系统
  • CMS教程
  • 游戏攻略
  • 脚本语言
  • 平面设计
  • 软件教程
  • 网络安全
  • 电脑知识
  • 服务器
  • 视频教程
  • windows
  • 服务器硬件
  • 服务器运维
  • 云计算
  • 虚拟化
  • IIS教程
  • Linux
  • Apache
  • Ftp
  • DNS
  • Nginx
您的位置:首页 > 服务器 >云计算 > Spark源码分析之worker节点启动driver和executor

Spark源码分析之worker节点启动driver和executor

作者:kaaosidao的博客 字体:[增加 减小] 来源:互联网

本文主要包含spark,源码分析,worker-driver等服务器相关知识,kaaosidao的博客希望可以进行参考

一、启动driver

 

1.首先在Master.scala类中执行schedule()方法,该方法主要有两个方法lanuchDriver()和launchExecutor()分别用来启动driver和executor。在master上面一旦可用资源发生变动或者有新的application提交进来之后就会调用该schedule()方法。

2.先去调度所有的driver,针对这些application采取严格的优先级.当该worker上面的可用内存大于application需要的内存,以及worker上面可用的core数目大于application中需要的core数目,通过launchDriver(worker, driver)才能够启动dirver,启动之后将该dirver就从正在等待的driver列表(ArrayBuffer)中删除同时将当前driver的状态设置为(launched=)true。

 

def launchDriver(worker: WorkerInfo, driver: DriverInfo) {
    logInfo("Launching driver " + driver.id + " on worker " + worker.id)
    worker.addDriver(driver)
    driver.worker = Some(worker)
    worker.actor ! LaunchDriver(driver.id, driver.desc)//向Worker节点发送一个基于AKKA ACTOR的事件通知模型的样例类LaunchDriver
    driver.state = DriverState.RUNNING
  }

 

3.worker接收到master发送来的LanuchDriver事件通知(Worker.scala类中)

 

case LaunchDriver(driverId, driverDesc) => {
      logInfo(s"Asked to launch driver $driverId")
      val driver = new DriverRunner(
        conf,
        driverId,
        workDir,
        sparkHome,
        driverDesc.copy(command = Worker.maybeUpdateSSLSettings(driverDesc.command, conf)),
        self,
        akkaUrl)
      drivers(driverId) = driver
      driver.start()

      coresUsed += driverDesc.cores
      memoryUsed += driverDesc.mem
    }
 

 

 

 

在LaunchRiver内部,创建DriverRunner对象(内部封装了一个线程),将DriverRunner对象添加进该Worker内部的dirver列表(HashMap<Key(driverID), Value(DriverRunner对象)>),调用DriverRunner对象的start方法,同时修改该worker进程的可用cores个数,内存个数。

4.查看driver.start()方法(DriverRunner.scala)

 

def start() = {
    new Thread("DriverRunner for " + driverId) {
      override def run() {
        try {
          val driverDir = createWorkingDirectory()
          val localJarFilename = downloadUserJar(driverDir)

          def substituteVariables(argument: String): String = argument match {
            case "{{WORKER_URL}}" => workerUrl
            case "{{USER_JAR}}" => localJarFilename
            case other => other
          }

          // TODO: If we add ability to submit multiple jars they should also be added here
          /*
            类似组织启动进程的命令如下:
            Storm jar jar-path classpath paramters
           */
          val builder = CommandUtils.buildProcessBuilder(driverDesc.command, driverDesc.mem,
            sparkHome.getAbsolutePath, substituteVariables)
          launchDriver(builder, driverDir, driverDesc.supervise)
        }
        catch {
          case e: Exception => finalException = Some(e)
        }
        //获取当前driver的启动状态(KILLED ERROR FINISHED FAILED)
        val state =
          if (killed) {
            DriverState.KILLED
          } else if (finalException.isDefined) {
            DriverState.ERROR
          } else {
            finalExitCode match {
              case Some(0) => DriverState.FINISHED
              case _ => DriverState.FAILED
            }
          }

        finalState = Some(state)

        worker ! DriverStateChanged(driverId, state, finalException)
      }
    }.start()
  }
该方法主要作用:①创建的driver工作的目录,②从master中下载要执行的jar包(移动计算、不移动数据),③调用内部的方法launchDriver(ProcessBuilder, driverDriver, supervise)来启动driver启动ProcessBuilder(使用java的API去操作一个java的进程|Process Runtime)。driver在此启动起来

 

并向worker发送DriverStateChanged(driverId, state, finalException)事件通知,worker接收到后,向master发送相同的DriverStateChanged(driverId, state, finalException)事件通知,

二、启动executor

1.类似于启动driver。在schedule()方法满足一定的条件(e.g:资源满足等等其他条件)后启动launchExecutor()  (Master.scala)

 

def launchExecutor(worker: WorkerInfo, exec: ExecutorDesc) {
    logInfo("Launching executor " + exec.fullId + " on worker " + worker.id)
    worker.addExecutor(exec)
    worker.actor ! LaunchExecutor(masterUrl,
      exec.application.id, exec.id, exec.application.desc, exec.cores, exec.memory)
    //启动executor之后,将该executor注册添加进driver
    exec.application.driver ! ExecutorAdded(
      exec.id, worker.id, worker.hostPort, exec.cores, exec.memory)
  }
2.在worker上面,创建了一个内部持有一个Thread线程的ExecutorRunner的对象,内部封装了启动该executor所需要的信息,调用start方法

 

 

case LaunchExecutor(masterUrl, appId, execId, appDesc, cores_, memory_) =>
......
    val manager = new ExecutorRunner(...)
    manager.start()
    coresUsed += cores_
    memoryUsed += memory_
    master ! ExecutorStateChanged(appId, execId, manager.state, None, None)
......
3. 二-2中manager.start()方法(ExecutorRunner.scala)

 

 

def start() {
    workerThread = new Thread("ExecutorRunner for " + fullId) {
      override def run() { fetchAndRunExecutor() }
    }
    workerThread.start()
    // Shutdown hook that kills actors on shutdown.
    shutdownHook = new Thread() {
      override def run() {
        killProcess(Some("Worker shutting down"))
      }
    }
    Runtime.getRuntime.addShutdownHook(shutdownHook)
  }
}

 
def fetchAndRunExecutor() { try { // Launch the process val builder = CommandUtils.buildProcessBuilder(appDesc.command, memory, sparkHome.getAbsolutePath, substituteVariables) val command = builder.command() logInfo("Launch command: " + command.mkString("\"", "\" \"", "\"")) builder.directory(executorDir) builder.environment.put("SPARK_LOCAL_DIRS", appLocalDirs.mkString(",")) // In case we are running this from within the Spark Shell, avoid creating a "scala" // parent process for the executor command builder.environment.put("SPARK_LAUNCH_WITH_SCALA", "0") // Add webUI log urls val baseUrl = s"http://$publicAddress:$webUiPort/logPage/?appId=$appId&executorId=$execId&logType=" builder.environment.put("SPARK_LOG_URL_STDERR", s"${baseUrl}stderr") builder.environment.put("SPARK_LOG_URL_STDOUT", s"${baseUrl}stdout") process = builder.start() val header = "Spark Executor Command: %s\n%s\n\n".format( command.mkString("\"", "\" \"", "\""), "=" * 40) // Redirect its stdout and stderr to files val stdout = new File(executorDir, "stdout") stdoutAppender = FileAppender(process.getInputStream, stdout, conf) val stderr = new File(executorDir, "stderr") Files.write(header, 
  


 
分享到:QQ空间新浪微博腾讯微博微信百度贴吧QQ好友复制网址打印

您可能想查找下面的文章:

  • spark core源码分析16 Shuffle详解-读流程,sparkshuffle
  • Machine Learning On Spark——第一节:基础数据结构(一),learningspark
  • spark core源码分析15 Shuffle详解-写流程,sparkshuffle
  • spark core源码分析14 参数配置,sparkcore
  • spark core源码分析10 Task的运行,sparkcore
  • spark core源码分析9 从简单例子看action操作,sparkcore
  • hadoop做HA后,hbase修改,hadoophahbase修改
  • Spark Core and Cluster Managers(翻译自Learning.Spark.Lightning-Fast.Big.Data.Analysis),
  • hadoop2.7配置HA,使用zk和journal,hadoop2.7zk
  • spark core源码分析7 Executor的运行,sparkexecutor

相关文章

  • 通过TelnetClient获取Zookeeper监控数据,zookeeperclient
  • 在streaming process中为什么需要类似sql查询语言,streamingprocess
  • 开源微内核seL4 microkernel,sel4microkernel
  • 程序员必备的代码审查(Code Review)清单,codereview
  • 【甘道夫】Spark1.3.0 Submitting Applications 官方文档精华摘要,spark1.3
  • 机器学习--决策树(ID3)算法案例,id3案例
  • 使用export/import导出和导入docker容器,exportdocker
  • **FASPOT 功能简介**,faspot功能简介
  • 《3》CentOS7.0+OpenStack+kvm云平台部署—配置Glance,《3》kvm
  • 修改 openstack 中 nova boot 创建实例只能在10个以内的限制,openstacknova

文章分类

  • windows
  • 服务器硬件
  • 服务器运维
  • 云计算
  • 虚拟化
  • IIS教程
  • Linux
  • Apache
  • Ftp
  • DNS
  • Nginx

最近更新的内容

    • Spark 批量写数据入HBase,spark数据入hbase
    • Slope One 协同过滤 推荐算法,slopeone
    • MapReduce中Shuffle过程整理,mapreduceshuffle
    • Spark源码分析之worker节点启动driver和executor
    • 机器学习-Python中训练模型的保存和再使用,-python模型
    • Ambari管理Hadoop集群时遇到的问题,ambarihadoop
    • 什么是Code Review,CodeReview
    • Hadoop之——自定义分组比较器实现分组功能,hadoop分组
    • hadoop权威指南(第四版)要点翻译(1)——Foreword and Preface,hadoopforeword
    • Error: Failed to launch instance &quot;win7&quot;: Please try again later [Error: No valid host was found. ].,win7valid

关于我们 - 联系我们 - 免责声明 - 网站地图

©2020-2025 All Rights Reserved. linkedu.com 版权所有