欧美一级特黄大片做受成人-亚洲成人一区二区电影-激情熟女一区二区三区-日韩专区欧美专区国产专区

好程序員大數(shù)據(jù)分享Spark任務(wù)和集群?jiǎn)?dòng)流程

好程序員大數(shù)據(jù)分享Spark任務(wù)和集群?jiǎn)?dòng)流程,Spark集群?jiǎn)?dòng)流程

從網(wǎng)站建設(shè)到定制行業(yè)解決方案,為提供成都網(wǎng)站設(shè)計(jì)、網(wǎng)站制作服務(wù)體系,各種行業(yè)企業(yè)客戶提供網(wǎng)站建設(shè)解決方案,助力業(yè)務(wù)快速發(fā)展。創(chuàng)新互聯(lián)建站將不斷加快創(chuàng)新步伐,提供優(yōu)質(zhì)的建站服務(wù)。

1.調(diào)用start-all.sh腳本,開始啟動(dòng)Master

2.Master啟動(dòng)以后,preStart方法調(diào)用了一個(gè)定時(shí)器,定時(shí)檢查超時(shí)的Worker后刪除

3.啟動(dòng)腳本會(huì)解析slaves配置文件,找到啟動(dòng)Worker的相應(yīng)節(jié)點(diǎn).開始啟動(dòng)Worker

4.Worker服務(wù)啟動(dòng)后開始調(diào)用preStart方法開始向所有的Master進(jìn)行注冊(cè)

5.Master接收到Worker發(fā)送過來的注冊(cè)信息,Master開始保存注冊(cè)信息并把自己的URL響應(yīng)給Worker

6.Worker接收到Master的URL后并更新,開始調(diào)用一個(gè)定時(shí)器,定時(shí)的向Master發(fā)送心跳信息

?

任務(wù)提交流程

1.Driver端會(huì)通過spark-submit腳本啟動(dòng)SaparkSubmit進(jìn)程,此時(shí)創(chuàng)建了一個(gè)非常重要的對(duì)象(SparkContext),開始向Master發(fā)送消息

2.Master接收到發(fā)送過來的信息后開始生成任務(wù)信息,并把任務(wù)信息放到一個(gè)對(duì)列里

3.Master把所有有效的Worker過濾出來,按照空閑的資源進(jìn)行排序

4.Master開始向有效的Worker通知拿取任務(wù)信息并啟動(dòng)相應(yīng)的Executor

5.Worker啟動(dòng)Executor并向Driver反向注冊(cè)

6.Driver開始把生成的task發(fā)送給相應(yīng)的Executor,Executor開始執(zhí)行任務(wù)

?

集群?jiǎn)?dòng)流程

1.首先創(chuàng)建Master類

import akka.actor.{Actor, ActorSystem, Props}

import com.typesafe.config.{Config, ConfigFactory}

?

import scala.collection.mutable

import scala.concurrent.duration._

?

class Master(val masterHost: String, val masterPort: Int) extends Actor{

?

??// 用來存儲(chǔ)Worker的注冊(cè)信息

??val idToWorker = new mutable.HashMap[String, WorkerInfo]()

?

??// 用來存儲(chǔ)Worker的信息

??val workers = new mutable.HashSet[WorkerInfo]()

?

??// Worker的超時(shí)時(shí)間間隔

??val checkInterval: Long = 15000

?

?

??// 生命周期方法,在構(gòu)造器之后,receive方法之前只調(diào)用一次

??override def preStart(): Unit = {

????// 啟動(dòng)一個(gè)定時(shí)器,用來定時(shí)檢查超時(shí)的Worker

????import context.dispatcher

????context.system.scheduler.schedule(0 millis, checkInterval millis, self, CheckTimeOutWorker)

??}

?

??// 在preStart方法之后,不斷的重復(fù)調(diào)用

??override def receive: Receive = {

????// Worker -> Master

????case RegisterWorker(id, host, port, memory, cores) => {

??????if (!idToWorker.contains(id)){

????????val workerInfo = new WorkerInfo(id, host, port, memory, cores)

????????idToWorker += (id -> workerInfo)

????????workers += workerInfo

?

????????println("a worker registered")

?

????????sender ! RegisteredWorker(s"akka.tcp://${Master.MASTER_SYSTEM}" +

??????????s"@${masterHost}:${masterPort}/user/${Master.MASTER_ACTOR}")

??????}

????}

????case HeartBeat(workerId) => {

??????// 通過傳過來的workerId獲取對(duì)應(yīng)的WorkerInfo

??????val workerInfo: WorkerInfo = idToWorker(workerId)

??????// 獲取當(dāng)前時(shí)間

??????val currentTime = System.currentTimeMillis()

??????// 更新最后一次心跳時(shí)間

??????workerInfo.lastHeartbeatTime = currentTime

????}

????case CheckTimeOutWorker => {

??????val currentTime = System.currentTimeMillis()

??????val toRemove: mutable.HashSet[WorkerInfo] =

????????workers.filter(w => currentTime - w.lastHeartbeatTime > checkInterval)

?

??????// 將超時(shí)的Worker從idToWorker和workers中移除

??????toRemove.foreach(deadWorker => {

????????idToWorker -= deadWorker.id

????????workers -= deadWorker

??????})

?

??????println(s"num of workers: ${workers.size}")

????}

??}

}

object Master{

??val MASTER_SYSTEM = "MasterSystem"

??val MASTER_ACTOR = "Master"

?

??def main(args: Array[String]): Unit = {

????val host = args(0)

????val port = args(1).toInt

?

????val configStr =

??????s"""

?????????|akka.actor.provider = "akka.remote.RemoteActorRefProvider"

?????????|akka.remote.netty.tcp.hostname = "$host"

?????????|akka.remote.netty.tcp.port = "$port"

??????""".stripMargin

?

????// 配置創(chuàng)建Actor需要的配置信息

????val config: Config = ConfigFactory.parseString(configStr)

?

????// 創(chuàng)建ActorSystem

????val actorSystem: ActorSystem = ActorSystem(MASTER_SYSTEM, config)

?

????// 用actorSystem實(shí)例創(chuàng)建Actor

????actorSystem.actorOf(Props(new Master(host, port)), MASTER_ACTOR)

?

????actorSystem.awaitTermination()

?

??}

}

2.創(chuàng)建RemoteMsg特質(zhì)

trait RemoteMsg extends Serializable{

?

}

?

// Master -> self(Master)

case object CheckTimeOutWorker

?

// Worker -> Master

case class RegisterWorker(id: String, host: String,

??????????????????????????port: Int, memory: Int, cores: Int) extends RemoteMsg

?

// Master -> Worker

case class RegisteredWorker(masterUrl: String) extends RemoteMsg

?

// Worker -> self

case object SendHeartBeat

?

// Worker -> Master(HeartBeat)

case class HeartBeat(workerId: String) extends RemoteMsg

3.創(chuàng)建Worker類

import java.util.UUID

?

import akka.actor.{Actor, ActorRef, ActorSelection, ActorSystem, Props}

import com.typesafe.config.{Config, ConfigFactory}

?

import scala.concurrent.duration._

?

class Worker(val host: String, val port: Int, val masterHost: String,

?????????????val masterPort: Int, val memory: Int, val cores: Int) extends Actor{

?

??// 生成一個(gè)Worker ID

??val workerId = UUID.randomUUID().toString

?

??// 用來存儲(chǔ)MasterURL

??var masterUrl: String = _

?

??// 心跳時(shí)間間隔

??val heartBeat_interval: Long = 10000

?

??// master的Actor

??var master: ActorSelection = _

?

??override def preStart(){

????// 獲取Master的Actor

????master = context.actorSelection(s"akka.tcp://${Master.MASTER_SYSTEM}" +

??????s"@${masterHost}:${masterPort}/user/${Master.MASTER_ACTOR}")

?

????master ! RegisterWorker(workerId, host, port, memory, cores)

??}

?

??override def receive: Receive = {

????// Worker接收到Master發(fā)送過來的注冊(cè)成功的信息(masterUrl)

????case RegisteredWorker(masterUrl) => {

??????this.masterUrl = masterUrl

??????// 啟動(dòng)一個(gè)定時(shí)器,定時(shí)給Master發(fā)送心跳

??????import context.dispatcher

??????context.system.scheduler.schedule(0 millis, heartBeat_interval millis, self, SendHeartBeat)

????}

????case SendHeartBeat => {

??????// 向Master發(fā)送心跳

??????master ! HeartBeat(workerId)

????}

?

??}

?

}

object Worker{

??val WORKER_SYSTEM = "WorkerSystem"

??val WORKER_ACTOR = "Worker"

?

??def main(args: Array[String]): Unit = {

????val host = args(0)

????val port = args(1).toInt

????val masterHost = args(2)

????val masterPort = args(3).toInt

????val memory = args(4).toInt

????val cores = args(5).toInt

?

????val configStr =

??????s"""

?????????|akka.actor.provider = "akka.remote.RemoteActorRefProvider"

?????????|akka.remote.netty.tcp.hostname = "$host"

?????????|akka.remote.netty.tcp.port = "$port"

??????""".stripMargin

?

????// 配置創(chuàng)建Actor需要的配置信息

????val config: Config = ConfigFactory.parseString(configStr)

?

????// 創(chuàng)建ActorSystem

????val actorSystem: ActorSystem = ActorSystem(WORKER_SYSTEM, config)

?

????// 用actorSystem實(shí)例創(chuàng)建Actor

????val worker: ActorRef = actorSystem.actorOf(

??????Props(new Worker(host, port, masterHost, masterPort, memory, cores)), WORKER_ACTOR)

?

????actorSystem.awaitTermination()

?

??}

}

4.創(chuàng)建初始化類

class WorkerInfo(val id: String, val host: String, val port: Int,

?????????????????val memory: Int, val cores: Int) {

?

??// 初始化最后一次心跳的時(shí)間

??var lastHeartbeatTime: Long = _

?

}

5.本地測(cè)試需要傳入?yún)?shù):

好程序員大數(shù)據(jù)分享Spark任務(wù)和集群?jiǎn)?dòng)流程

網(wǎng)站題目:好程序員大數(shù)據(jù)分享Spark任務(wù)和集群?jiǎn)?dòng)流程
當(dāng)前地址:http://www.aaarwkj.com/article26/gjcecg.html

成都網(wǎng)站建設(shè)公司_創(chuàng)新互聯(lián),為您提供移動(dòng)網(wǎng)站建設(shè)、建站公司、網(wǎng)站制作網(wǎng)站營銷、網(wǎng)站改版微信小程序

廣告

聲明:本網(wǎng)站發(fā)布的內(nèi)容(圖片、視頻和文字)以用戶投稿、用戶轉(zhuǎn)載內(nèi)容為主,如果涉及侵權(quán)請(qǐng)盡快告知,我們將會(huì)在第一時(shí)間刪除。文章觀點(diǎn)不代表本網(wǎng)站立場(chǎng),如需處理請(qǐng)聯(lián)系客服。電話:028-86922220;郵箱:631063699@qq.com。內(nèi)容未經(jīng)允許不得轉(zhuǎn)載,或轉(zhuǎn)載時(shí)需注明來源: 創(chuàng)新互聯(lián)

微信小程序開發(fā)
人妻中文字幕av资源| 亚洲精品日韩在线欧美| 久久精品色妇熟妇丰满人妻| 亚洲国产色一区二区三区| 日本少妇人妻一区二区| 欧美香蕉一区二区视频| 婷婷国产综合一区二区三区| 日本黄色中文字幕在线观看| 日本区一区二区三视频| 一本色道久久亚洲综合精品蜜桃| 国产精品亚洲在线视频| 国产成人精品亚洲日本片| 丁香色婷婷国产精品视频| 中午字幕久久亚洲精品| 欧美日韩亚洲人人夜夜澡| 亚洲一区二区在线视频在线观看| 国产精品一区日韩专区| 亚洲日本高清一二三区| 国产精品黄色片在线观看| 亚洲精品在线观看午夜福利| 日本高清不卡在线播放| 高潮国产精品一区二区| 欧美美女福利午夜视频| 色综合视频二区偷拍在线| 亚洲福利一区二区三区| 欧美伊人久久大综合精品| 国产福利精品一区二区av| 日本av在线中文一区二区| 很黄很刺激的视频中文字幕| 欧美日韩国产天堂一区| 日韩亚洲欧美成人一区| 日韩精品一二三黄色一级| av电影网站中文字幕| 亚洲综合日韩精品在线| 精品日韩av高清一区二区三区| 九九热这里只有免费视频| 亚洲人成伊人久久成| 熟妇丰满多毛的大阴户| 国产亚洲综合久久系列| 操老熟女一区二区三区| 欧美精品在线高清观看|