大数据工程项目开发实战活页式教程 课件 第6章 离线处理辅助系统_第1页
大数据工程项目开发实战活页式教程 课件 第6章 离线处理辅助系统_第2页
大数据工程项目开发实战活页式教程 课件 第6章 离线处理辅助系统_第3页
大数据工程项目开发实战活页式教程 课件 第6章 离线处理辅助系统_第4页
大数据工程项目开发实战活页式教程 课件 第6章 离线处理辅助系统_第5页
已阅读5页,还剩50页未读 继续免费阅读

下载本文档

版权说明:本文档由用户提供并上传,收益归属内容提供方,若内容存在侵权,请进行举报或认领

文档简介

1第6章

离线处理辅助系统目录01Spark概述03SparkSQL05Spark城市旅游热力图02SparkCore04SparkStreamingSpark概述

01Spark概述Spark,快速、通用、可扩展的大数据分析引擎,支持多种计算模式。Spark优点Spark之所以能够成为目前最为流行的内存计算框架:快、易用、通用、兼容性Spark特点Spark提供了RDD、SQL和Streaming等多种功能,SparkRDD用来做离线数据分析的。SparkSQL是进行SQL分析的。SparkStreaming是做实时计算的。Spark的功能Spark概述Spark生态系统,包括SparkCore、Streaming、SQL、GraphX和MLib等组件。Spark生态系统Spark概述上传Spark压缩包至Master,解压、重命名、配置环境变量与文件。Spark环境部署SparkCore

02RDD转换操作RDD转换是延迟执行的,只有遇到动作才实际运行,提升Spark效率。RDD依赖RDD依赖分为窄依赖和宽依赖,前者类似独生子女,后者类似超生。RDD概述RDD是Spark中的核心数据抽象,支持容错和高效并行计算。RDDActionAction操作会直接触发RDD计算,包括reduce、reduceByKey等算子。RDD缓存RDD缓存使Spark能快速重用内存中的数据集,加速迭代算法和交互查询。SparkCore什么是RDDSparkCoreRDD(ResilientDistributedDataset)叫做分布式数据集,是Spark中最基本的数据抽象,它代表一个不可变、可分区、里面的元素可并行计算的集合。RDD具有数据流模型的特点:自动容错、位置感知性调度和可伸缩性。RDD允许用户在执行多个查询时显式地将工作集缓存在内存中,后续的查询能够重用工作集,这极大地提升了查询速度。RDD的属性SparkCore1)一组分片(Partition),即数据集的基本组成单位。2)一个计算每个分区的函数。3)RDD之间的依赖关系。4)一个Partitioner,即RDD的分片函数。5)一个列表,存储存取每个Partition的优先位置(preferredlocation)。为什么会产生RDDSparkCore传统的MapReduce虽然具有自动容错、平衡负载和可拓展性的优点,但是其最大缺点是采用非循环式的数据流模型,使得在迭代计算式要进行大量的磁盘IO操作。RDD正是解决这一缺点的抽象方法

RDD是Spark提供的最重要的抽象的概念,它是一种有容错机制的特殊集合,可以分布在集群的节点上,以函数式编操作集合的方式,进行各种并行操作。可以将RDD理解为一个具有容错机制的特殊集合,它提供了一种只读、只能有已存在的RDD变换而来的共享内存,然后将所有数据都加载到内存中,方便进行多次重用。为什么会产生RDDSparkCoreRDD的容错机制实现分布式数据集容错方法有两种:数据检查点和记录更新RDD。采用记录更新的方式:记录所有更新点的成本很高。所以,RDD只支持粗颗粒变换,即只记录单个块上执行的单个操作,然后创建某个RDD的变换序列(血统)存储下来;变换序列指,每个RDD都包含了他是如何由其他RDD变换过来的以及如何重建某一块数据的信息。因此RDD的容错机制又称“血统”容错。要实现这种“血统”容错机制,最大的难题就

是如何表达父RDD和子RDD之间的依赖关系。实际上依赖关系可以分两种,窄依赖和宽依赖:窄依赖:子RDD中的每个数据块只依赖于父RDD 中对应的有限个固定的数据块;宽依赖:子RDD中的一个数据块可以依赖于父RDD中的所有数据块。为什么会产生RDDSparkCoreRDD内部的设计每个RDD都需要包含以下四个部分:①

源数据分割后的数据块,源代码中的splits变量。②

关于“血统”的信息,源码中的dependencies变量。③

一个计算函数(该RDD如何通过父RDD计算得到),源码中的iterator(split)和compute函数。④

一些关于如何分块和数据存放位置的元信息,如源码中的partitioner和preferredLocations。RDD在Spark中的地位及作用SparkCore为什么会有Spark?因为传统的并行计算模型无法有效的解决迭代计算(iterative)和交互式计算(interactive);而Spark的使命便是解决这两个问题,这也是他存在的价值和理由。Spark如何解决迭代计算?其主要实现思想就是RDD,把所有计算的数据保存在分布式的内存中。迭代计算通常情况下都是对同一个数据集做反复的迭代计算,数据在内存中将大大提升IO操作。这也是Spark涉及的核心:内存计算。RDD在Spark中的地位及作用SparkCoreSpark如何实现交互式计算?因为Spark是用scala语言实现的,Spark和scala能够紧密的集成,所以Spark可以完美的运用scala的解释器,使得其中的scala可以向操作本地集合对象一样轻松操作分布式数据集。

Spark和RDD的关系?可以理解为:RDD是一种具有容错性基于内存的集群计算抽象方法,Spark则是这个抽象方法的实现。如何操作RDD?SparkCore(1)由一个已经存在的Scala集合创建。valrdd1=sc.parallelize(Array(1,2,3,4,5,6,7,8))(2)由外部存储系统的数据集创建,包括本地的文件系统,还有所有Hadoop支持的数据集,比如HDFS、Cassandra、HBase等

valrdd2=sc.textFile("hdfs://master:9000/words.txt")通过上面两种方式,我们可以很容易的创建出RDD,从而完成Spark的内存计算,那么,关于RDD算子的使用,有两种关键的方式分别为Transformation和Action,接下来我们详细讲解下这两种操作是如何转换RDD的。RDD转换操作RDD转换是延迟执行的,只有遇到动作才实际运行,提升Spark效率。RDD依赖RDD依赖分为窄依赖和宽依赖,前者类似独生子女,后者类似超生。RDD概述RDD是Spark中的核心数据抽象,支持容错和高效并行计算。RDDActionAction操作会直接触发RDD计算,包括reduce、reduceByKey等算子。RDD缓存RDD缓存使Spark能快速重用内存中的数据集,加速迭代算法和交互查询。SparkCoreSparkCoreSparkCoreRDD转换操作RDD转换是延迟执行的,只有遇到动作才实际运行,提升Spark效率。RDD依赖RDD依赖分为窄依赖和宽依赖,前者类似独生子女,后者类似超生。RDD概述RDD是Spark中的核心数据抽象,支持容错和高效并行计算。RDDActionAction操作会直接触发RDD计算,包括reduce、reduceByKey等算子。RDD缓存RDD缓存使Spark能快速重用内存中的数据集,加速迭代算法和交互查询。SparkCoreSparkCoreRDD转换操作RDD转换是延迟执行的,只有遇到动作才实际运行,提升Spark效率。RDD依赖RDD依赖分为窄依赖和宽依赖,前者类似独生子女,后者类似超生。RDD概述RDD是Spark中的核心数据抽象,支持容错和高效并行计算。RDDActionAction操作会直接触发RDD计算,包括reduce、reduceByKey等算子。RDD缓存RDD缓存使Spark能快速重用内存中的数据集,加速迭代算法和交互查询。SparkCoreSparkCoreRDD转换操作RDD转换是延迟执行的,只有遇到动作才实际运行,提升Spark效率。RDD依赖RDD依赖分为窄依赖和宽依赖,前者类似独生子女,后者类似超生。RDD概述RDD是Spark中的核心数据抽象,支持容错和高效并行计算。RDDActionAction操作会直接触发RDD计算,包括reduce、reduceByKey等算子。RDD缓存RDD缓存使Spark能快速重用内存中的数据集,加速迭代算法和交互查询。SparkCoreSparkCoreSpark速度非常快的原因之一,就是在不同操作中可以在内存中持久化或缓存个数据集。当持久化某个RDD后,每一个节点都将把计算的分片结果保存在内存中,并在对此RDD或衍生出的RDD进行的其他动作中重用。这使得后续的动作变得更加迅速。RDD相关的持久化和缓存,是Spark最重要的特征之一。可以说,缓存是Spark构建迭代式算法和快速交互式查询的关键。SparkCoreDAG在Spark中根据RDD间的依赖关系划分Stage,实现分布式内存计算。Spark运行架构Checkpoint是Spark中用于可靠持久化中间计算结果,提高容错性和性能的功能。checkpointSpark编程通过RDD实现数据处理,支持WordCount等任务。SparkRDD编程基础SparkCore是Spark的核心组件,帮助完成内存计算,但其代码书写方式对SQL工程师不友好。总结SparkCoreDAG(DirectedAcyclicGraph)叫做有向无环图,原始的RDD通过一系列的转换就就形成了DAG,根据RDD之间的依赖关系的不同将DAG划分成不同的Stage,对于窄依赖,partition的转换处理在Stage中完成计算。对于宽依赖,由于有Shuffle的存在,只能在parentRDD处理完成后,才能开始接下来的计算,因此宽依赖是划分Stage的依据SparkCoreDAG在Spark中的应用:SparkCoreSpark的任务提交机制:SparkCoreDAG在Spark中根据RDD间的依赖关系划分Stage,实现分布式内存计算。Spark运行架构Checkpoint是Spark中用于可靠持久化中间计算结果,提高容错性和性能的功能。checkpointSpark编程通过RDD实现数据处理,支持WordCount等任务。SparkRDD编程基础SparkCore是Spark的核心组件,帮助完成内存计算,但其代码书写方式对SQL工程师不友好。总结SparkCorecheckpoint原理机制:当RDD使用cache机制从内存中读取数据,如果数据没有读到,会使用checkpoint机制读取数据。此时如果没有checkpoint机制,那么就需要找到父RDD重新计算数据了,因此checkpoint是个很重要的容错机制。checkpoint就是对于一个RDDchain(链),如果中间某些中间结果RDD,后面需要反复使用该数据,可能因为一些故障导致该中间数据丢失,那么就可以针对该RDD启动checkpoint机制。SparkCoreDAG在Spark中根据RDD间的依赖关系划分Stage,实现分布式内存计算。Spark运行架构Checkpoint是Spark中用于可靠持久化中间计算结果,提高容错性和性能的功能。checkpointSpark编程通过RDD实现数据处理,支持WordCount等任务。SparkRDD编程基础SparkCore是Spark的核心组件,帮助完成内存计算,但其代码书写方式对SQL工程师不友好。总结SparkCoreDAG在Spark中根据RDD间的依赖关系划分Stage,实现分布式内存计算。Spark运行架构Checkpoint是Spark中用于可靠持久化中间计算结果,提高容错性和性能的功能。checkpointSpark编程通过RDD实现数据处理,支持WordCount等任务。SparkRDD编程基础SparkCore是Spark的核心组件,帮助完成内存计算,但其代码书写方式对SQL工程师不友好。总结SparkSQL

03SparkSQL提升执行效率,兼容Hive,支持多种数据源访问。学习SparkSQL原因SparkSQL是Spark处理结构化数据的模块,提供DataFrame抽象和分布式SQL查询功能。什么是SparkSQLSparkSQL概述DataFrame是分布式数据容器,支持结构化和嵌套数据类型,提供友好API。DataFramesDataFrame常用操作包括DSL和SQL风格的查询、过滤和分组等。DataFrame常用操作SparkSQL可方便操作Hive表,实现数据分析与处理。SparkSQL与Hive结合SparkSQL编程编写SparkSQL查询程序,实现数据处理与分析功能。01编程执行SparkSQL查询SparkSQL通过JDBC实现与MySQL的数据交互,包括读取和写入操作。02MySQL数据源SparkSQL外部数据源操作掌握DataFrame操作、Hive数据处理及编程方式,减少代码量,支持数据库工程师进行大数据离线分析。SparkSQL学习下节转入SparkStreaming流式计算,探讨其与离线计算的不同之处。即将学习总结SparkStreaming

04SparkStreaming简介SparkStreaming与Storm均为优秀流式计算框架,主要区别在于语言和编程模型。SparkStorm对比分析SparkStreaming:高吞吐量、强容错的流式数据处理工具,支持多种数据源。SparkStreaming简介Spark与Storm的对比:SparkStreaming什么是DStreamDStream是SparkStreaming的基础抽象,代表持续性数据流及操作结果。DStream相关操作DStreams输出操作OutputOperations使DStream数据可输出至外部系统,触发实际计算。DStream操作分转换和输出两类,含特殊原语如updateStateByKey。DStreams变换DStreams算子与RDD类似,特别介绍了UpdateStateByKey和WindowOperations。SparkStreaming核心概念DStream示意图:SparkStreamingDStream流程图:批处理:DStream算子:Window窗口:OutputOperationsonDStreams:SparkStreaming编程SparkStreaming实现网络文本实时词频统计,监听端口接收数据处理。实时WordCount实现掌握SparkStreaming编程,通过实践学习

温馨提示

  • 1. 本站所有资源如无特殊说明,都需要本地电脑安装OFFICE2007和PDF阅读器。图纸软件为CAD,CAXA,PROE,UG,SolidWorks等.压缩文件请下载最新的WinRAR软件解压。
  • 2. 本站的文档不包含任何第三方提供的附件图纸等,如果需要附件,请联系上传者。文件的所有权益归上传用户所有。
  • 3. 本站RAR压缩包中若带图纸,网页内容里面会有图纸预览,若没有图纸预览就没有图纸。
  • 4. 未经权益所有人同意不得将文件中的内容挪作商业或盈利用途。
  • 5. 人人文库网仅提供信息存储空间,仅对用户上传内容的表现方式做保护处理,对用户上传分享的文档内容本身不做任何修改或编辑,并不能对任何下载内容负责。
  • 6. 下载文件中如有侵权或不适当内容,请与我们联系,我们立即纠正。
  • 7. 本站不保证下载资源的准确性、安全性和完整性, 同时也不承担用户因使用这些下载资源对自己和他人造成任何形式的伤害或损失。

评论

0/150

提交评论