Spark Streaming源码解读之流数据不断接收详解

Spark Streaming源码解读之流数据不断接收详解

ID:37903287

大小:346.00 KB

页数:24页

时间:2019-06-02

Spark Streaming源码解读之流数据不断接收详解_第1页
Spark Streaming源码解读之流数据不断接收详解_第2页
Spark Streaming源码解读之流数据不断接收详解_第3页
Spark Streaming源码解读之流数据不断接收详解_第4页
Spark Streaming源码解读之流数据不断接收详解_第5页
资源描述:

《Spark Streaming源码解读之流数据不断接收详解》由会员上传分享,免费在线阅读,更多相关内容在教育资源-天天文库

1、SparkStreaming源码解读之流数据不断接收详解特别说明:在上一遍文章中有详细的叙述Receiver启动的过程,如果不清楚的朋友,请您查看上一篇博客,这里我们就基于上篇的结论,继续往下说。博文的目标是:SparkStreaming在接收数据的全生命周期贯通组织思路如下:a)接收数据的架构模式的设计b)然后再具体源码分析接收数据的架构模式的设计1.当有SparkStreaming有application的时候SparkStreaming会持续不断的接收数据。2.一般Receiver和Driver不在一个进程中的,所以接收到数据之后

2、要不断的汇报给Driver。3.SparkStreaming要接收数据肯定要使用消息循环器,循环器不断的接收到数据之后,然后将数据存储起来,再将存储完的数据汇报给Driver。4.SparkStreaming数据接收的过程也是MVC的架构,M是model也就是Receiver.C是Control也就是存储级别的ReceiverSupervisor。V是界面。5.ReceiverSupervisor是控制器,Receiver的启动是靠ReceiverTracker启动的,Receiver接收到数据之后是靠ReceiverSuperviso

3、r存储数据的。然后Driver就获得元数据也就是界面,通过界面来操作底层的数据,这个元数据就相当于指针。SparkStreaming接收数据流程如下:具体源码分析1.ReceiverTracker通过发送Job的方式,并且每个Job只有一个Task,并且Task中只通过一个ReceiverSupervisor启动一个Receiver.2.下图就是Receiver启动的流程图,现在就从ReceiverTracker的start开始今天的旅程。3.Start方法中创建Endpoint实例/**Starttheendpointandrecei

4、verexecutionthread.*/defstart():Unit=synchronized{if(isTrackerStarted){thrownewSparkException("ReceiverTrackeralreadystarted")}if(!receiverInputStreams.isEmpty){endpoint=ssc.env.rpcEnv.setupEndpoint("ReceiverTracker",newReceiverTrackerEndpoint(ssc.env.rpcEnv))if(!skipRec

5、eiverLaunch)launchReceivers()logInfo("ReceiverTrackerstarted")trackerState=Started}}4.LaunchReceivers源码如下:/***GetthereceiversfromtheReceiverInputDStreams,distributesthemtothe*workernodesasaparallelcollection,andrunsthem.*/privatedeflaunchReceivers():Unit={valreceivers=re

6、ceiverInputStreams.map(nis=>{valrcvr=nis.getReceiver()rcvr.setReceiverId(nis.id)rcvr})runDummySparkJob()logInfo("Starting"+receivers.length+"receivers")//此时的endpoint就是前面实例化的ReceiverTrackerEndpointendpoint.send(StartAllReceivers(receivers))}5.从图上可以知道,send发送消息之后,ReceiverTr

7、ackerEndpoint的receive就接收到了消息。overridedefreceive:PartialFunction[Any,Unit]={//LocalmessagescaseStartAllReceivers(receivers)=>valscheduledLocations=schedulingPolicy.scheduleReceivers(receivers,getExecutors)for(receiver<-receivers){valexecutors=scheduledLocations(receiver.s

8、treamId)updateReceiverScheduledExecutors(receiver.streamId,executors)receiverPreferredLocations(receive

当前文档最多预览五页,下载文档查看全文

此文档下载收益归作者所有

当前文档最多预览五页,下载文档查看全文
温馨提示:
1. 部分包含数学公式或PPT动画的文件,查看预览时可能会显示错乱或异常,文件下载后无此问题,请放心下载。
2. 本文档由用户上传,版权归属用户,天天文库负责整理代发布。如果您对本文档版权有争议请及时联系客服。
3. 下载前请仔细阅读文档内容,确认文档内容符合您的需求后进行下载,若出现内容与标题不符可向本站投诉处理。
4. 下载文档时可能由于网络波动等原因无法下载或下载错误,付费完成后未能成功下载的用户请联系客服处理。