(版本定制)第10课:Spark Streaming源码解读
admin
2023-01-31 12:07:16
0

本期内容:

    1、数据接收架构设计模式

    2、数据接收源码彻底研究


1、Receiver接受数据的过程类似于MVC模式:

Receiver,ReceiverSupervisor和Driver的关系相当于Model,Control,View,也就是MVC。

Model就是Receiver,存储数据Control,就是ReceiverSupervisor,Driver是获得元数据,也就是View。

(版本定制)第10课:Spark Streaming源码解读

2、数据的位置信息会被封装到RDD里面。

3、Receiver接受数据,交给ReceiverSupervisor去存储数据。

4、ReceiverTracker是通过发送一个又一个的Job,每个Job只有一个Task,每个Task里面就只有一个ReceiverSupervisor,用这个函数启动每一个Receiver。


下面我们简单的看下Receiver启动流程,应用程序首先通过JobScheduler的start方法来启动receiverTracker的start方法:

def start(): Unit = synchronized {
if (eventLoop != null) return // scheduler has already been started

logDebug("Starting JobScheduler")
eventLoop = new EventLoop[JobSchedulerEvent]("JobScheduler") {
override protected def onReceive(event: JobSchedulerEvent): Unit = processEvent(event)

override protected def onError(e: Throwable): Unit = reportError("Error in job scheduler", e)
  }
eventLoop.start()

// attach rate controllers of input streams to receive batch completion updates
for {
    inputDStream <- ssc.graph.getInputStreams
    rateController <- inputDStream.rateController
} ssc.addStreamingListener(rateController)

listenerBus.start(ssc.sparkContext)
receiverTracker = new ReceiverTracker(ssc)
inputInfoTracker = new InputInfoTracker(ssc)
receiverTracker.start() //receiver启动
jobGenerator.start()
  logInfo("Started JobScheduler")
}

通过调用receiverTracker.start()方法来进行一系列的操作:

/** Start the endpoint and receiver execution thread. */
def start(): Unit = synchronized {
if (isTrackerStarted) {
throw new SparkException("ReceiverTracker already started")
  }

if (!receiverInputStreams.isEmpty) {
endpoint = ssc.env.rpcEnv.setupEndpoint(
"ReceiverTracker", new ReceiverTrackerEndpoint(ssc.env.rpcEnv)) //Rpc消息通信,获取receiver的状态
if (!skipReceiverLaunch) launchReceivers() //启动receiver
    logInfo("ReceiverTracker started")
trackerState = Started
}
}

下面通过画图简单的描述下Receiver启动的内部机制:

(版本定制)第10课:Spark Streaming源码解读


参考博客:http://blog.csdn.net/hanburgud/article/details/51471047

                 http://lqding.blog.51cto.com/9123978/1774426

相关内容

热门资讯

我们的“文化体力”是如何被耗尽... 《花束般的恋爱》2026年已经过半,要盘点前半年爆火的流行语,“赛博确诊”一定是其中之一。大批年轻人...
跨国药企正在疯狂买中国创新药,... ► 文 观察者网心智观察所2026年的中国创新药行业,正在发生一个耐人寻味的场景。一边,全球最大的跨...
知名装机up主被曝负债200万... 据都市快报消息:7月21日,深耕高端PC圈的B站UP主(B站、小红书、抖音等平台上传视频的创作者、博...
AI脸让人生理性厌恶?“恐怖谷... 【文/观察者网 夏曼波】最近,“对AI脸感到不适”相关话题频繁上热搜,网友纷纷表示“不适”、“生理性...
要管管国会山股神了?美国参议院... 作者 | 第一财经 樊志菁当地时间周三,美国议会众议院以232票对198票表决结果,通过一项名为《停...
济南一广场突现“知了大军”密密... 近日,有济南市民在社交媒体发布视频称:在当地彩石附近的一处广场地面上密密麻麻聚集了很多类似知了的昆虫...
伊朗军队:打击科威特境内多处美... 当地时间23日,伊朗军队在一份公告中指出,几小时前,伊朗军队对科威特多哈营地的美军弹药和后勤仓库、阿...
洛阳“珠宝大王”年永安沉浮:涉... 澎湃新闻记者 王鑫 实习生 林晨 林宜之河南洛阳的“珠宝大王”年永安淡出公众视野多年后,这个曾把“金...
聚焦核心AI项目,亚马逊裁撤其... 亚马逊在本周三宣布新动态,公司正在裁撤通用人工智能(AGI)部门的部分员工,这标志着亚马逊在持续向人...
上海交大发布第七届“十大科技进... 7月22日,上海交通大学第七届“十大科技进展”在闵行大零号湾成果转化中心发布。本届入选成果包括6项基...