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

好程序員大數(shù)據(jù)分享Spark任務(wù)和集群啟動流程-創(chuàng)新互聯(lián)

好程序員大數(shù)據(jù)分享Spark任務(wù)和集群啟動流程,Spark集群啟動流程

創(chuàng)新新互聯(lián),憑借10余年的成都做網(wǎng)站、網(wǎng)站設(shè)計(jì)經(jīng)驗(yàn),本著真心·誠心服務(wù)的企業(yè)理念服務(wù)于成都中小企業(yè)設(shè)計(jì)網(wǎng)站有千余家案例。做網(wǎng)站建設(shè),選創(chuàng)新互聯(lián)

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

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

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

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

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

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

任務(wù)提交流程

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

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

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

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

5.Worker啟動Executor并向Driver反向注冊

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

集群啟動流程

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{

// 用來存儲Worker的注冊信息

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

// 用來存儲Worker的信息

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

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

val checkInterval: Long = 15000

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

override def preStart(): Unit = {

// 啟動一個(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獲取對應(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

// 用來存儲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ā)送過來的注冊成功的信息(masterUrl)

case RegisteredWorker(masterUrl) => {

this.masterUrl = masterUrl

// 啟動一個(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.本地測試需要傳入?yún)?shù):

好程序員大數(shù)據(jù)分享Spark任務(wù)和集群啟動流程

另外有需要云服務(wù)器可以了解下創(chuàng)新互聯(lián)scvps.cn,海內(nèi)外云服務(wù)器15元起步,三天無理由+7*72小時(shí)售后在線,公司持有idc許可證,提供“云服務(wù)器、裸金屬服務(wù)器、高防服務(wù)器、香港服務(wù)器、美國服務(wù)器、虛擬主機(jī)、免備案服務(wù)器”等云主機(jī)租用服務(wù)以及企業(yè)上云的綜合解決方案,具有“安全穩(wěn)定、簡單易用、服務(wù)可用性高、性價(jià)比高”等特點(diǎn)與優(yōu)勢,專為企業(yè)上云打造定制,能夠滿足用戶豐富、多元化的應(yīng)用場景需求。

當(dāng)前文章:好程序員大數(shù)據(jù)分享Spark任務(wù)和集群啟動流程-創(chuàng)新互聯(lián)
分享地址:http://aaarwkj.com/article8/idcip.html

成都網(wǎng)站建設(shè)公司_創(chuàng)新互聯(lián),為您提供定制開發(fā)、定制網(wǎng)站網(wǎng)站設(shè)計(jì)、品牌網(wǎng)站建設(shè)、微信公眾號、網(wǎng)站維護(hù)

廣告

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

商城網(wǎng)站建設(shè)
日韩欧美一区二区三区不卡在线| 我要看亚洲黄色片一级| 熟女精品国产一区二区三区| 中文字幕人妻熟女人妻| 日本欧美国产一区二区| 久久人妻蜜桃一区二区三区| 久久婷亚洲综合五月天| 日本美女阴部毛茸茸视频| 激情少妇一区二区三区| 国产高清大片一级黄色| 五月婷婷六月丁香伊人网| 亚洲精品一区二区日本| 九九在线视频免费观看精品视频| 欧美黄片完整版在线观看| 日本精品视频一区二区三区| 亚洲一区二区三区精品国产| 五月婷婷六月丁香综合激情| 久久精品国产一区二区| 中文字幕日韩精品在线看| 和富婆啪啪一区二区免费看| 色桃子av一区二区三区| 欧美在线观看黄片视频| av中文字幕亚洲一区二区| 免费人成黄页网站在线播放国产| 999热这里只有精品视频| 欧美欧成人一区二区三区a∨| 亚洲伊人av第一页在线观看| 日韩亚洲国产欧美在线观看| 91精品蜜臀国产综合久久久久久| 国产在线精品91系列| 国产传媒在线视频免费| 97视频观看免费观看| 最新国产毛片久热精品视频| 四虎精品国产一区二区三区| 亚洲国产视频不卡一区| 九色视频在线观看91| 精品午夜人妻一区二区| 未满十八禁止观看免费观看| 欧美精品一区二区毛卡片| 国产原创传媒在线观看| 亚洲国产av福利久久|