欢迎来到天天文库
浏览记录
ID:37903287
大小:346.00 KB
页数:24页
时间:2019-06-02
《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
此文档下载收益归作者所有