Spark大数据技术 14.Spark Core编程之创建RDD_第1页
Spark大数据技术 14.Spark Core编程之创建RDD_第2页
Spark大数据技术 14.Spark Core编程之创建RDD_第3页
Spark大数据技术 14.Spark Core编程之创建RDD_第4页
Spark大数据技术 14.Spark Core编程之创建RDD_第5页
已阅读5页,还剩16页未读 继续免费阅读

下载本文档

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

文档简介

SparkCore编程之创建RDD深入解析RDD的两种创建方式及其应用Catalogue目录1.SparkRDD基础回顾重温RDD弹性分布式数据集定义,解析SparkContext核心入口与交互式Shell环境配置。2.内存集合创建RDD详解parallelize与makeRDD创建方法,深入剖析分区机制对计算并行度的影响。3.外部存储创建RDD掌握textFile读取本地与HDFS文件的方法,理解文件读取的默认分区规则与参数调优。4.创建方式对比分析对比内存与外部创建的异同,明确不同数据源场景下的RDD创建选型与性能考量。5.实战演练与答疑通过代码实操巩固RDD创建流程,解析常见报错原因,总结最佳实践与优化技巧。分布式(Distributed)数据被切分为多个分区并分布式存储在集群不同节点上,实现并行处理;作为不可变的元素集合,可容纳各类数据对象,为大规模数据的高效计算提供了分布式基础架构支持。弹性(Resilient)拥有血缘关系(Lineage)记录演变过程,分区数据丢失时可重新计算恢复,无需全量备份;支持根据计算需求动态调整分区数量,灵活适配集群资源,实现高容错与弹性伸缩。RDD:Spark的核心抽象RDD:Spark计算的起点RDD(弹性分布式数据集)是Spark的核心数据抽象,代表一个不可变、可分区且支持并行计算的元素集合。在Spark中,一切计算都基于RDD展开,数据的转换与结果的输出均需通过RDD完成,它是构建分布式计算逻辑的基础与起点。01.创建(Creation)从内存集合(如List)或外部存储系统(HDFS、文件)构建初始RDD。这是分布式计算的起点,决定了数据的来源与初始形态,是后续所有操作的基础。02.转换(Transformation)通过map、filter等算子将现有RDD转换为新的RDD。该过程采用懒加载机制,仅记录依赖关系(血缘),不会立即触发集群上的实际计算任务。03.行动(Action)触发集群上的真正计算,将结果返回给Driver程序或写入外部存储(如count、collect、saveAsTextFile)。只有执行行动操作,之前的所有转换逻辑才会被调度执行。SparkContext:应用与集群的桥梁SparkContext核心功能与交互式环境特性解析01核心入口定位——SparkContext(简称sc)是Spark应用的总控入口,负责与集群建立连接,协调Driver与Executor之间的通信,是分布式计算的基础枢纽。02构建分布式数据集——提供parallelize、makeRDD、textFile等核心方法,将本地数据或外部存储文件转化为弹性分布式数据集(RDD),开启分布式计算流程。03管理共享变量资源——支持创建广播变量(Broadcast)实现大数据集高效分发,以及累加器(Accumulator)实现分布式场景下的计数与统计,优化集群资源利用。04集群资源动态协调——与集群管理器(YARN/Standalone/Mesos)通信,按需申请CPU、内存等计算资源,管理任务的分发与执行调度。05Shell环境自动初始化——在SparkShell交互式编程环境中,系统会自动创建并初始化SparkContext实例,直接赋值给sc变量,无需手动编写初始化代码即可快速上手。SparkShell:快速验证与学习的利器01交互式开发:即写即得的验证环境SparkShell是内置Scala/Python解释器的交互式命令行工具,支持逐行输入代码并实时返回执行结果,无需进行打包编译流程。它是初学者掌握SparkAPI的最佳入门环境,也是工程师快速验证数据处理逻辑、调试算法原型的高效工具,极大缩短了从想法到验证的迭代周期。02--master参数:灵活适配不同运行模式通过--master参数可指定执行环境:local为本地单线程模式,适合代码调试;local[n]启用n个线程并行计算;local[*]自动调用机器所有CPU核心(默认);spark://HOST:PORT用于连接独立Spark集群;yarn模式则对接HadoopYARN集群,支持弹性资源调度与大规模分布式计算,满足不同场景的算力需求。从内存集合创建RDDPART01核心方法:parallelize、makeRDD|适用场景:快速原型验证与小数据测试使用parallelize方法创建RDDparallelize是Spark初始化分布式数据集的基础方法,它将单机内存中的Scala本地集合(如List、Array)转化为弹性分布式数据集(RDD),实现数据从“单机存储”到“集群并行计算”的跨越,是快速构建测试数据或小规模数据处理的首选方式。01.数据源:seq(必选)传入Scala标准的序列集合(如List,Array),作为RDD的原始数据来源,支持任意数据类型的序列化集合。02.并行度:numSlices(可选)指定RDD的分区数量,默认等于集群CPU核心数。分区数直接决定了任务的并行执行粒度,需根据数据量合理设置。Scala代码实战:从本地列表构建RDDvaldata=List(1,2,3,4,5)//定义本地集合

valrdd=sc.parallelize(data,2)//并行化,设置2个分区

println(rdd.first())//触发行动算子,输出结果:1创建RDD:makeRDD方法▍函数签名(Scala)defmakeRDD[T](seq:Seq[T],numSlices:Int):RDD[T]•seq:待并行化的本地集合(数据源)

•numSlices:分区数量(可选,默认取集群默认并行度)▍快速上手示例//1.定义本地数组数据源

valwords=Array("Spark","Scala","Hadoop")

//2.并行化创建RDD

valrdd=sc.makeRDD(words)

//3.执行计算:统计元素个数

println(rdd.count())//输出结果:3▍高级特性:位置偏好makeRDD提供了特殊的重载方法,支持为数据指定位置偏好(PreferredLocations)。这允许开发者手动指定数据分区在集群节点上的存储或计算位置,最大化利用数据本地性(DataLocality),减少网络传输开销,是针对复杂场景进行性能调优的高级手段。分区的核心价值与优势决定计算并行度,分区数建议匹配或超过集群核心数;依托数据本地性将任务调度至数据所在节点,减少网络传输开销,同时支持灵活的分区策略,是Spark实现高效分布式计算的基石。什么是RDD分区?RDD的基本组成单元,将分布式数据集在逻辑上切分为多个独立的数据块,物理上分布于集群的不同节点。每个分区对应一个独立的计算任务(Task),是Spark实现数据并行计算的核心载体,支撑了分布式任务的并行调度与执行。深入理解RDD分区(Partition)Spark默认分区数的确定规则01本地模式:线程数直接映射在local或local[n]模式下,默认分区数等于指定的线程数。例如配置local[4]启动Spark时,RDD会默认创建4个分区,使分区数与本地并发线程数完全一致,最大化利用单机多核处理能力。02集群模式:取核心数与2的最大值在Standalone或YARN集群中,默认分区数为max(集群总可用核心数,2)。此规则确保分区数不会过少,保障分布式任务的基础并行度,避免因分区不足导致的资源闲置,是分布式计算性能保障的基础。📝代码验证:查看实际分区数//1.并行化创建RDD,不指定分区

valrdd=sc.parallelize(1to100)

//2.获取并打印默认分区数量

println("DefaultPartitions:"+rdd.getNumPartitions)

//结果:8核机器集群模式下输出8💡核心价值:理解默认分区规则是进行Spark性能调优的第一步,合理的分区数能最大化利用集群资源,避免数据倾斜。实战演示:从内存创建RDD01启动SparkShell执行命令进入交互式环境:spark-shell--masterlocal[2]以本地2核模式启动,自动初始化SparkContext(sc),适合快速开发调试。02并行化构建RDD基于本地集合创建分布式数据集:sc.parallelize(dataList,3)将本地列表转为分布式RDD,显式指定3个分区,决定任务的并行计算粒度。03透视分区数据遍历分区并打印内容:mapPartitionsWithIndex(func).collect()通过转换算子关联分区索引,行动算子触发计算,直观验证数据分布。//核心代码实现(Scala)//1.准备本地数据valnames=List("Alice","Bob","Charlie","David","Eve","Frank")//2.并行化创建RDD,指定3个分区valrdd=sc.parallelize(names,3)//3.打印各分区内容的函数定义与调用rdd.mapPartitionsWithIndex((i,it)=>{println(s"P$i:${it.mkString(",")}");Iterator.empty}).collect()➜终端执行与输出结果scala>rdd.getNumPartitions

res0:Int=3scala>rdd.mapPartitionsWithIndex(...).collect()

P0:Alice,David

P1:Bob,Eve

P2:Charlie,Frank

res1:Array[String]=Array()说明:数据被均匀打散到3个分区中,保证了后续并行计算的负载均衡。从外部存储创建RDDPART02核心方法:textFile与核心场景——它是Spark读取外部数据的核心入口,可高效加载本地文件、HDFS分布式文件系统及云存储(如S3)中的海量数据,适配结构化与非结构化数据的处理需求,是构建大数据分布式分析应用的基础步骤。RDD读取:textFile方法▍核心功能与定位Spark中读取外部文本数据的基础入口API,用于将本地文件系统或分布式存储(如HDFS)中的文本文件加载为RDD。它会将文件的每一行内容作为RDD的一个独立String类型元素,是构建大数据处理管道的基础步骤。▍Scala函数原型deftextFile(path:String,minPartitions:Int):RDD[String]•path:文件或目录的路径,支持file://(本地)和hdfs://(分布式)协议;

•minPartitions:可选参数,指定RDD的最小分区数量,决定了数据并行处理的粒度下限。▍底层工作机制方法按换行符(\n)分割文件内容,每行文本独立成为RDD的一个分区元素。若路径指向目录,则递归加载目录下所有文件;支持对大文件自动分片,对小文件进行合并优化,确保数据分区的均衡性。▍核心特性总结1.惰性求值:仅记录数据依赖,触发Action算子才实际读取;

2.广泛兼容:适用于TXT、CSV、日志文件等各类文本格式;

3.弹性分区:可通过分区数参数灵活控制并行度,优化性能。01本地文件直读——使用file://前缀指定本地路径(如file:///data/file.txt),可直接读取本地单机文件系统中的数据,适合本地调试场景。02HDFS分布式读取——前缀为hdfs://<namenode>:<port>,也可简写为绝对路径(如/user/data/),无缝适配Hadoop分布式文件系统的海量数据存储。03目录整体加载——直接传入目录路径(如sc.textFile("/log_dir/")),Spark会自动递归读取该目录下的所有文件,无需逐个指定文件名。04通配符模糊匹配——支持*、?等通配符筛选特定模式文件,例如sc.textFile("/data/*.log")可批量读取所有日志后缀的文件,灵活匹配多源数据。⚠️集群运行关键注意——若在集群模式下使用本地路径(file://),必须确保该文件存在于集群中每一个Worker节点的相同物理路径下,否则会抛出FileNotFoundException异常。掌握Spark读取文件的四大路径模式与集群规则SparktextFile灵活的路径访问机制textFile的分区策略01核心规则:块大小决定分区数Spark会为文件的每个HDFS块创建一个对应的RDD分区,HDFS默认块大小为128MB。这是分区的核心依据,单个文件的物理存储块与RDD分区形成一一映射,直接决定了计算任务的基础并行度,是理解Spark数据读取并行化的关键前提。02minPartitions:分区数的下限保障minPartitions是分区数量的下限阈值,并非决定值。若按块大小计算的分区数小于该值,Spark会自动将分区数扩容至该下限,避免因分区过少导致并行度不足。例如300MB文件按128MB/块分为3个块(对应3个分区),若minPartitions设为2则保持3个,若设为4则会扩容至4个分区。实战演示:读取本地wc.txt文件01准备工作:创建测试数据源在节点的/opt/spark/目录下创建wc.txt文件,内容如下:SparkScalaJava

KafkaFlinkHadoop

HDFSMapReduceSpark02核心操作:读取本地文件生成RDD使用SparkContext的textFile方法,需指定file://协议://加载本地文件

vallocalRDD=sc.textFile("file:///opt/spark/wc.txt")▶查看分区:getNumPartitionsscala>localRDD.getNumPartitions

res0:Int=1说明:因文件尺寸远小于HDFS块大小(128M),默认仅创建1个分区。▶收集数据:collect()scala>localRDD.collect()

res1:Array[String]=Array(...)//返回所有行数组注意:collect()会将分布式数据拉取到Driver端,仅用于小数据集调试。实战演示:读取HDFS文件并进行单词计数01.前置准备:上传文件至HDFS分布式存储在执行计算前,需将本地数据源文件上传至HDFS,确保集群节点可访问。使用如下Shell命令完成上传:执行命令:hdfsdfs-put/opt/spark/wc.txt/spark/(将本地wc.txt上传至HDFS的/spark目录)02.Spark单词计数核心逻辑(ScalaRDD算子)通过转换算子与行动算子的配合,实现从文件读取到结果输出的全流程://1.读取文件->2.扁平化拆分->3.键值映射->4.聚合统计valcounts=sc.textFile("/spark/wc.txt").flatMap(_.split("")).map((_,1)).reduceByKey(_+_)counts.collect()//触发计算,输出:Array((Spark,2),(Hadoop,1),(Scala,3))💡核心要点:textFile支持读取HDFS路径,reduceByKey利用Shuffle机制实现分布式环境下的高效聚合。”从外部存储创建(textFile/分布式文件)数据来源与规模:支持本地文件、HDFS、S3等分布式存储系统,专为TB级以上的大规模离线数据设计,是生产环境的标准选择。核心应用场景:企业级ETL数据清洗、海量日志分析、历史数据批处理等对数据量要求极

温馨提示

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

评论

0/150

提交评论