首页 文章 精选 留言 我的

精选列表

搜索[编译原理],共10000篇文章
优秀的个人博客,低调大师

Spark Streaming 原理剖析

通过源码呈现 Spark Streaming 的底层机制。 1. 初始化与接收数据 Spark Streaming 通过分布在各个节点上的接收器缓存接收到的流数据并将流数 据 包 装 成 Spark 能 够 处 理 的 RDD的格式 输入到Spark Streaming 之 后由Spark Streaming将作业提交到Spark集群进行执行如图1所示。 图 1 Spark Streaming 执行模型 初始化的过程主要可以概括为两点 1调度器的初始化。 调度器调度 Spark Streaming 的运行用户可以通过配置相关参数进行调优。 2将输入流的接收器转化为 RDD 在集群进行分布式分配然后启动接收器集合中的每个接收器。针对不同的数据源 Spark Streaming 提供了不同的数据接收器分布在各个节点上的每个接收器可以认为是一个特定的进程接收一部分流数据作为输入。 用户也可以针对自身生产环境状况自定义开发相应的数据接收器。 如图 2 所示接收器分布在各个节点上。通过下面代码创建并行的、在不同Worker 节点分布的 receiver 集合。 val tempRDD = if (hasLocationPreferences) { val receiversWithPreferences = receivers.map(r => (r, Seq(r.preferredLocation.get))) ssc.sc.makeRDD[Receiver[_]](receiversWithPreferences) } else { // 在这里创造 RDD 相当于进入 SparkContext.makeRDD // 此处将 receivers 的集合作为一个 RDD 进行分区 RDD[Receiver] // 即使是只有一个输入流按照这个分布式也是流的输入端在 worker 而不再 Master … // 将 receivers 的集合打散然后启动它们 … ssc.sparkContext.runJob(tempRDD, startReceiver) … } 图 2 Spark Streaming 接收器 2. 数据接收与转化 在上面的“初始化与接收数据”部分中已经介绍过 receiver 集合转换为 RDD在集群上分布式地接收数据流。那么每个 receiver 是怎样接收并处理数据流的呢读者可以通过图 3对输入流的处理有一个全面的了解。图 3为 Spark Streaming 数据接收与转化的示意图。 图 3 的主要流程如下。 1数据缓冲在 receiver 的 receive 函数中接收流数据将接收到的数据源源不断地放入到 BlockGenerator.currentBuffer。 2缓冲数据转化为数据块在 BlockGenerator 中有一个定时器RecurringTimer将 当 前 缓 冲 区 中 的 数 据 以 用 户 定 义 的 时 间 间 隔 封 装 为 一 个 数 据 块 Block 放 入 到 BlockGenerator 的 blocksForPush 队列中这个队列。 3数据块转化为 Spark 数据块在 BlockGenerator 中有一个 BlockPushingThread线程不断地将 blocksForPush 队列中的块传递给 BlockManager让 BlockManager 将 数据存储为块。 BlockManager 负责 Spark 中的块管理。 4元数据存储在 pushArrayBuffer 方法中还会将已经由 BlockManager 存储的元数据信息例如 Block 的 id 号传递给 ReceiverTracker ReceiverTracker 会将存储的 blockId 放到对应 StreamId 的队列中。 图 3 Spark Streaming 数据接收与转化 图中部分组件的作用如下 KeepPushingBlocks调用此方法持续写入和保持数据块。 pushArrayBuffer调用 pushArrayBuffer 方法将数据块存储到 BlockManager 中。 reportPushedBlock存储完成后汇报数据块信息到主节点。 receivedBlockInfo Meta Data已经接收到的数据块元数据记录。 streamId数据流 Id。 BlockInfo数据块元数据信息。 BlockManager.put数据块存储器写入备份数据块到其他节点。 Receiver数据块接收器接收数据块。 BlockGenerator数据块生成器将数据缓存生成 Spark 能处理的数据块。 BlockGenerator.currentBuffer缓存网络接收的数据记录等待之后转换为 Spark的数据块。 BlockGenerator.blocksForPushing将一块连续数据记录暂存为数据块待后续转换为 Spark 能够处理的 BlockManager 中的数据块A Block As a BlockManager’sBlock)。 BlockGenerator.blockPushingThread守护线程负责将数据块转换为 BlockManager中数据块。 ReceiveTracker输入数据块的元数据管理器负责管理和记录数据块。 BlockManager Spark 数据块管理器负责数据块在内存或磁盘的管理。 RecurringTimer时间触发器每隔一定时间进行缓存数据的转换。 上面的过程中涉及最多的类就是 BlockGenerator在数据转化的过程中其扮演者不可或缺的角色。 private[streaming] class BlockGenerator( listener: BlockGeneratorListener, receiverId: Int, conf: SparkConf ) extends Logging 3. 生成 RDD 与提交 Spark Job Spark Streaming 根据时间段将数据切分为 RDD然后触发 RDD 的 Action 提交 Job Job 被 提 交 到 Job Manager 中 的 Job Queue 中 由 Job Scheduler 调 度 之 后Job Scheduler 将 Job 提交到 Spark 的 Job 调度器然后将 Job 转换为大量的任务分发给 Spark 集群执行如图 4 所示。 图4 Spark Streaming 调度模型 Job generator 中通过下面的方法生成 Job 进行调度和执行。 从下面的代码可以看出 job 是从 outputStream 中生成的然后再触发反向回溯执行 整个 DStream DAG类似 RDD 的机制。 private def generateJobs(time: Time) { SparkEnv.set(ssc.env) Try(graph.generateJobs(time)) match { case Success(jobs) => // 获取输入数据块的元数据信息 val receivedBlockInfo = graph.getReceiverInputStreams.map { stream => . . . }.toMap jobScheduler.submitJobSet(JobSet(time, jobs, receivedBlockInfo)) case Failure(e) => jobScheduler.reportError("Error generating jobs for time " + time, e) } eventActor !DoCheckpoint(time) } // 下 面 进 入 JobScheduler 的 submitJobSet 方 法 一 探 究 竟 JobScheduler 是 整 个 Spark Streaming 调度的核心组件 def submitJobSet(jobSet: JobSet) { . . . jobSets.put(jobSet.time, jobSet) jobSet.jobs.foreach(job => jobExecutor.execute(new JobHandler(job))) . . . } // 进入 Graph 生成 job 的方法 Graph 本质是 DStreamGraph 类生成的对象 final private[streaming] class DStreamGraph extends Serializable with Logging { def generateJobs(time: Time): Seq[Job] = { . . . private val inputStreams = new ArrayBuffer[InputDStream[_]]() private val outputStreams = new ArrayBuffer[DStream[_]]() . . . val jobs = this.synchronized { outputStreams.flatMap(outputStream => outputStream.generateJob(time)) . . . } // outputStreams 中的对象是 DStream下面进入 DStream 的 generateJob 一探究竟 private[streaming] def generateJob(time: Time): Option[Job] = { getOrCompute(time) match { case Some(rdd) => { val jobFunc = () => { val emptyFunc = { (iterator: Iterator[T]) => {} } // 此处相当于针对每个时间段生成的一个 RDD会调用 SparkContext 的方法 runJob 提交 Spark 的一 个 Job context.sparkContext.runJob(rdd, emptyFunc) } Some(new Job(time, jobFunc)) } case None => None } } // 在 DStream 算是父类一些具体的 DStream 例如 SocketInputStream 等的类的父类可以通过 SocketInputDStream 看是如何通过上面的 getOrCompute 生成 RDD 的 private[streaming] def getOrCompute(time: Time): Option[RDD[T]] = { generatedRDDs.get(time) match { . . . case None => { if (isTimeValid(time)) { // Dstream 是个父类这里代表的是子类的 compute 方法 DStream 通过 compute 调用用户自定 义函数。当任务执行时同一个 stage 中的 DStream 函数会串联依次执行 compute(time) match { . . . generatedRDDs.put(time, newRDD) . . . } 在 SocketInputDStream 的 compute 方法中生成了对应时间片的 RDD override def compute(validTime: Time): Option[RDD[T]] = { if (validTime >= graph.startTime) { val blockInfo = ssc.scheduler.receiverTracker.getReceivedBlockInfo(id) receivedBlockInfo(validTime) = blockInfo val blockIds = blockInfo.map(_.blockId.asInstanceOf[BlockId]) Some(new BlockRDD[T](ssc.sc, blockIds)) } else { Some(new BlockRDD[T](ssc.sc, Array[BlockId]())) } } Spark Streaming 在保证实时处理的要求下还能够保证高吞吐与容错性。用户的数据分析中很多情况下也存在需要分析图数据运行图算法通过 GraphX 可以简便地开发分布式图分析算法。 本文转自大数据躺过的坑博客园博客原文链接http://www.cnblogs.com/zlslch/p/5725374.html如需转载请自行联系原作者

优秀的个人博客,低调大师

LeakCanary内存检测原理

LeakCanary介绍 LeakCanary内存检测工具是由squar公司开源的著名项目。此项目主要用于内存检测。开源鲜明的目录结构如下: leakcanary | |-leakcanary-analyzer |-leakcanary-android |-leakcanary-android-no-op |-leakcanary-watcher |-leakcanary-simple leakcanary-android,leakcanary-android-no-op是针对android的检查封装。leakcanary-analyzer:用于分析dump文件leakcanary-watcher:用于监控内存泄露leakcanary-a

优秀的个人博客,低调大师

Hadoop数据读写原理

数据流 MapReduce作业(job)是客户端执行的单位:它包括输入数据、MapReduce程序和配置信息。Hadoop把输入数据划分成等长的小数据发送到MapReduce,称之为输入分片。Hadoop为每个分片创建一个map任务,由它来运行用户自定义的map函数来分析每个分片中的记录。 这里分片的大小,如果分片太小,那么管理分片的总时间和map任务创建的总时间将决定作业的执行的总时间。对于大数据作业来说,一个理想的分片大小往往是一个HDFS块的大小,默认是64MB(可以通过配置文件指定) map任务的执行节点和输入数据的存储节点是同一节点时,Hadoop的性能达到最佳。这就是为什么最佳分片的大小与块大小相同,它是最大的可保证存储在单个节点上的数据量如果分区跨越两个块,那么对于任何一个HDFS节点而言,基本不可能同时存储着两数据块,因此此分布的某部分必须通过网络传输到节点,这与使用本地数据运行map任务相比,显然效率很低。 reduce任务并不具备数据本地读取的优势,一个单一的reduce的任务的输入往往来自于所有mapper的输出。因此,有序map的输出必须通过网络传输到reduce任务运行的节点,并在哪里进行合并,然后传递到用户自定义的reduce函数中。 一般情况下,多个reduce任务的数据流成为"shuffle",因为每个reduce任务的输入都由许多map任务来提供。 Hadoop流 流适用于文字处理,在文本模式下使用时,它有一个面向行的数据视图。map的输入数据把标准输入流传输到map函数,其中是一行一行的传输,然后再把行写入标准输出。该框架调用mapper的map()方法来处理读入的每条记录,然而map程序可以决定如何处理输入流,可以轻松地读取和同一时间处理多行,用户的java map实现是压栈记录,但它仍可以考虑处理多行,具体做法是将mapper中实例变量中之前的行汇聚在一起(可用其他语言实现)。 HDFS的设计 HDFS是为以流式数据访问模式存储超大文件而设计的文件系统,在商用硬件的集群上运行。 流式数据访问:一次写入、多次读取模式是最高效的,一个数据集通常由数据源生成或复制,接着在此基础上进行各种各样的分析。 低延迟数据访问:需要低延迟访问数据在毫秒范围内的应用不适用于HDFS,HDFS是为达到高数据吞吐量而优化的,这有可能会以延迟为代价。(低延迟访问可以参考HBASE) 大量的小文件:namenode存储着文件系统的元数据,文件数量的限制也由namenode的内存量决定。每个文件,索引目录以及块占大约150个字节,因此,如果有一百万文件,每个文件占一个块,就至少需要300MB的内存。 多用户写入,任意修改文件:HDFS中的文件只有一个写入者。 HDFS的块比磁盘的块大,目的是为了减少寻址的开销。通过让一个块足够大,从磁盘转移数据的时间能够远远大于定位这个开始端的时间。因此,传送一个由多个块组成的文件的时间就取决于磁盘传送率。 文件读取与写入 HDFS中读取数据 客户端是通过调用fileSystem对象的open()来读取希望打开的文件的。对于HDFS,这个对象是分布式文件系统的一个实例。 (1)DistributedFileSystem通过使用RPC来调用namenode,以确定文件开头部分的块的位置,对于每一个块,namenode返回具有该块副本的数据节点地址。随后这些数据节点根据它们与客户端的距离来排序,如果该客户端本身就是一个数据节点,便从本地数据节点读取。(Distributed FileSystem返回一个FSData InputStream转而包装了一个DFSInputStream对象) (2)存储着文件开头部分的块的数据节点地址的DFSInputStream随机与这些块的最近的数据节点相连接,通过在数据流中重复调用read(),数据就会从数据节点返回客户端。到达块的末端时,DFSInputSteam会关闭与数据节点间的连接,然后为下一个块找到最佳的数据节点。 (3)客户端从流中读取数据时,块是按照DFSInputStream打开与数据节点的新连接的顺序读取的。它也会调用namenode来检索下一组需要的块的数据节点的位置。一旦客户端完成读取,就对文件系统数据输入调用close()。 这个设计的重点是,客户端直接联系数据节点去检索数据,通过namenode指引到每个块中最好的数据节点。因为数据流动在此集群中是在所有数据节点分散进行的,因此这种设计能使HDFS可扩展到最大的并发客户端数量。namenode提供块位置请求,其数据是存储在内存,非常的高效。 文件写入 客户端是通过在DistributedFilesystem中调用create()来创建文件,DistributedFilesystem一个RPC去调用namenode,在文件系统的命名空间中创建一个新文件。namenode执行各种不同的检查以确保这个文件不会已经存在,并且客户端有创建文件的权限。文件系统数据输出流控制一个DFSoutPutstream,负责datanode与namenode之间的通信。 客户端完成数据的写入后,会在流中调用clouse(),在向namenode发送完信息之前,此方法会将余下的所有包放入datanode管线并等待确认,namenode节点已经知道文件由哪些块组成(通过Data streamer询问块分配),所以它值需在返回成功前等待块进行最小量的复制。 一旦写入的数据超过一个块的数据,新的读取者就能看见第一个块,对于之后的块也是这样。总之,它始终是当前正在被写入的块,其他读取者是看不见的。HDFS提供一个方法来强制所有的缓存与datanode同步,即在文件系统数据输出流调用sync()方法,在syno()返回成功后,HDFS能保证文件中直至写入的最后的数据对所有新的读取者而言,都是可见且一致的。万一发生冲突(与客户端或HDFS),也不会造成数据丢失。 应用设计的重要性:如果不调用sync(),那么一旦客户端或系统发生故障,就可能失去一个块的数据,应该在适当的地方调用sync(). 通过distcp进行并行复制:Hadoop有一个叫distcp(分布式复制)的有用程序,能从Hadoop的文件系统并行复制大量数据。如果集群在Hadoop的同一版本上运行,就适合使用hdfs方案: hadoop distcp hdfs://namenode1/foo hdfs://namenode2/bar 将从第一个集群中复制/foo目录到第二个集群中的/bar目录下 参考:《Hadoop权威指南-第四版》

资源下载

更多资源
Mario

Mario

马里奥是站在游戏界顶峰的超人气多面角色。马里奥靠吃蘑菇成长,特征是大鼻子、头戴帽子、身穿背带裤,还留着胡子。与他的双胞胎兄弟路易基一起,长年担任任天堂的招牌角色。

Nacos

Nacos

Nacos /nɑ:kəʊs/ 是 Dynamic Naming and Configuration Service 的首字母简称,一个易于构建 AI Agent 应用的动态服务发现、配置管理和AI智能体管理平台。Nacos 致力于帮助您发现、配置和管理微服务及AI智能体应用。Nacos 提供了一组简单易用的特性集,帮助您快速实现动态服务发现、服务配置、服务元数据、流量管理。Nacos 帮助您更敏捷和容易地构建、交付和管理微服务平台。

Spring

Spring

Spring框架(Spring Framework)是由Rod Johnson于2002年提出的开源Java企业级应用框架,旨在通过使用JavaBean替代传统EJB实现方式降低企业级编程开发的复杂性。该框架基于简单性、可测试性和松耦合性设计理念,提供核心容器、应用上下文、数据访问集成等模块,支持整合Hibernate、Struts等第三方框架,其适用范围不仅限于服务器端开发,绝大多数Java应用均可从中受益。

Sublime Text

Sublime Text

Sublime Text具有漂亮的用户界面和强大的功能,例如代码缩略图,Python的插件,代码段等。还可自定义键绑定,菜单和工具栏。Sublime Text 的主要功能包括:拼写检查,书签,完整的 Python API , Goto 功能,即时项目切换,多选择,多窗口等等。Sublime Text 是一个跨平台的编辑器,同时支持Windows、Linux、Mac OS X等操作系统。

用户登录
用户注册