ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

头歌实践教学平台:Spark大数据编程(三十六~四十)

2026/8/17 18:22:27 拓冰建站 浏览量
头歌实践教学平台:Spark大数据编程(三十六~四十) 三十六、RDD概述任务描述本关任务根据下面的相关知识完成与数据认知相关的选择题。相关知识RDD介绍RDD 是Spark的核心抽象即 弹性分布式数据集residenta distributed dataset。代表一个不可变可分区里面元素可并行计算的集合。其具有数据流模型的特点自动容错位置感知性调度和可伸缩性。在Spark中对数据的所有操作不外乎创建RDD、转化已有RDD以及调用 RDD操作进行求值。RDD结构图RDD具有五大特性1、一组分片Partition即数据集的基本组成单位RDD是由一系列的partition组成的。将数据加载为RDD时一般会遵循数据的本地性一般一个HDFS里的block会加载为一个partition。2、RDD之间的依赖关系。依赖还具体分为宽依赖和窄依赖但并不是所有的RDD都有依赖。为了容错重算cachecheckpoint也就是说在内存中的RDD操作时出错或丢失会进行重算。3、由一个函数计算每一个分片。Spark中的RDD的计算是以分片为单位的每个RDD都会实现compute函数以达到这个目的。compute函数会对迭代器进行复合不需要保存每次计算的结果。4、可选如果RDD里面存的数据是key-value形式则可以传递一个自定义的Partitioner进行重新分区。5、可选RDD提供一系列最佳的计算位置即数据的本地性。RDD之间的依赖关系RDD之间有一系列的依赖关系依赖关系又分为窄依赖和宽依赖。窄依赖父RDD和子RDD partition之间的关系是一对一的。或者父RDD一个partition只对应一个子RDD的partition情况下的父RDD和子RDD partition关系是多对一的也可以理解为没有触发shuffle。宽依赖父RDD与子RDD partition之间的关系是一对多。父RDD的一个分区的数据去到子RDD的不同分区里面。也可以理解为触发了shuffle。特别说明对于join操作有两种情况如果join操作的使用每个partition仅仅和已知的Partition进行join此时的join操作就是窄依赖其他情况的join操作就是宽依赖。RDD创建1、从Hadoop文件系统或与Hadoop兼容的其他持久化存储系统如Hive、Cassandra、HBase输入例如HDFS创建。2、通过集合进行创建。算子算子可以分为Transformation 转换算子和Action 行动算子。RDD是懒执行的如果没有行动操作出现所有的转换操作都不会执行。编程要求根据相关知识按照要求完成右侧选择题任务包含单选题和多选题。三十七、算子概述三十八、Spark RDD编程初级实践第1关数据去重任务描述本关任务编写Spark独立应用程序实现数据去重。相关知识为了完成本关任务你需要掌握RDD的创建RDD的转换操作RDD的行动操作。RDD的创建使用textFile()方法从本地文件系统中加载数据创建RDD示例如下val lines sc.textFile(file:///home/hadoop/word.txt)执行sc.textFile()方法以后Spark从本地文件word.txt中加载数据到内存在内存中生成一个RDD对象lines这个RDD里面包含了若干个元素每个元素的类型是String类型也就是说从word.txt文件中读取出来的每一行文本内容都成为RDD中的一个元素。使用map()函数转换得到相应的键值对RDD示例如下val lines sc.textFile(file:///home/hadoop/word.txt)val pairRDD lines.flatMap(line line.split( )).map(word (word,1))上面示例中map(word(word,1))函数的作用是取出RDD中的每个元素也就是每个单词赋值给word然后把word转换成 (word,1) 的键值对形式。RDD的转换操作对于RDD而言每一次转换操作都会产生新的RDD供给下一个操作使用。RDD的转换过程是惰性求值的也就是说整个转换过程只是记录了转换的轨迹并不会发生真正的计算只有遇到行动操作时才会触发真正的计算。常见的RDD转换操作如下所示**filter(func)**筛选出满足函数func的元素并返回一个新的RDD示例如下val lines sc.textFile(file:///home/hadoop/word.txt)val linesWithSpark lines.filter(line line.contains(Spark))**map(func)**将每个元素传递到函数func中并将结果返回为一个新的RDD示例如下val lines sc.textFile(file:///home/hadoop/word.txt)val words lines.map(line line.split( ))**groupByKey()**应用于(K,V)键值对RDD时返回一个新的(K, Iterable)形式的RDD示例如下val lines sc.textFile(file:///home/hadoop/word.txt)val words lines.flatMap(line line.split( )).map(word (word,1)).groupByKey()**sortByKey()**应用于(K,V)键值对RDD时返回一个新的根据key排序的RDD示例如下val pairRDD sc.parallelize(Array(Hadoop,3),(Spark,5),(Hive,2))pairRDD.sortByKey().foreach(println)输出(Hadoop,3)(Hive,2)(Spark,5)**partitionBy(partitioner: Partitioner)**根据partitioner函数生成新的ShuffleRDD将原RDD重新分区示例如下val lines sc.textFile(file:///home/hadoop/word.txt, 3)val words lines.flatMap(line line.split( )).map(word (word,1)).partitionBy(new HashPartitioner(1))keys将键值对RDD中所有元素的key返回形成一个新的RDD示例如下val pairRDD sc.parallelize(Array(Hadoop,3),(Spark,5),(Hive,2))pairRDD.keys.foreach(println)输出HadoopSparkHiveRDD的行动操作对于RDD而言只有遇到行动操作时才会执行“从头到尾”的真正的计算从文件中加载数据完成一次又一次转换操作最终完成行动操作得到结果。常见的RDD行动操作如下所示**count()**返回RDD中元素的个数**collect()**以数组的形式返回RDD中的所有元素**first()**返回RDD中的第一个元素**take(n)**以数组的形式返回RDD中的前n个元素**reduce(func)**通过函数func(输入两个参数并返回一个值)聚合RDD中的元素**foreach(func)**将RDD中的每个元素传递到函数func中运行 下面通过一个示例来介绍上述行动操作如下所示val rdd sc.parallelize(Array(1,2,3,4,5))println(rdd.count)println(rdd.first)println(rdd.take(3).mkString(, ))println(rdd.collect().mkString(- ))rdd.foreach(println)输出 5 1 1, 2, 3 1- 2- 3- 4- 5 1 2 3 4 5import org.apache.spark.SparkContextimport org.apache.spark.SparkConfimport org.apache.spark.HashPartitionerobject RemDup {def main(args: Array[String]): Unit {val conf new SparkConf().setAppName(RemDup).setMaster(local)val sc new SparkContext(conf)//输入文件fileA.txt和fileB.txt已保存在本地文件系统/root/step1_files目录中val dataFile file:///root/step1_filesval data sc.textFile(dataFile, 2)/********** Begin **********///第一步执行过滤操作把空行丢弃。val filteredData data.filter(line line.trim.nonEmpty)//第二步执行map操作取出RDD中每个元素去除尾部空格并生成一个(key, value)键值对。val pairRDD filteredData.map { line val parts line.trim.split(\\s) // 按任意空白符分割(parts(0) parts(1), parts(0) parts(1)) // 键值对都用完整行方便后续去重和排序}//第三步执行groupByKey操作把所有key相同的value都组织成一个value-list去重核心val groupedRDD pairRDD.groupByKey()//第四步对RDD进行重新分区变成一个分区保证排序全局有序val partitionedRDD groupedRDD.partitionBy(new HashPartitioner(1))//第五步执行sortByKey操作对RDD中所有元素都按照key的升序排序。val sortedRDD partitionedRDD.sortByKey(true) // true表示升序//第六步执行keys操作将键值对RDD中所有元素的key返回形成一个新的RDDkey就是去重后的行val resultRDD sortedRDD.keys//第七步执行collect操作以数组的形式返回RDD中所有元素。val resultArray resultRDD.collect()//第八步执行foreach操作并使用println打印出数组中每个元素的值。println() // 保留原有打印行resultArray.foreach(println)/********** End **********/}}第2关整合排序任务描述本关任务编写Spark独立应用程序实现整合排序。相关知识为了完成本关任务你需要掌握RDD的创建RDD的转换操作RDD的行动操作。编程要求假设某个目录下有多个文本文件每个文件中每一行内容均为一个整数。要求读取所有文件中的整数进行排序后输出到一个新的文件中输出的内容为每行两个整数第一个整数为第二个整数的排序位次第二个整数为原待排序的整数。import org.apache.spark.SparkContextimport org.apache.spark.SparkConfimport org.apache.spark.HashPartitionerobject FileSort {def main(args: Array[String]): Unit {val conf new SparkConf().setAppName(FileSort).setMaster(local)val sc new SparkContext(conf)//输入文件file1.txt、file2.txt和file3.txt已保存在本地文件系统/root/step2_files目录中val dataFile file:///root/step2_filesval data sc.textFile(dataFile, 3)/********** Begin **********///第一步执行过滤操作把空行丢弃。val filteredData data.filter(line line.trim.nonEmpty)//第二步执行map操作取出RDD中每个元素去除尾部空格并转换成整数生成一个(key, value)键值对。val pairRDD filteredData.map { line val num line.trim.toInt // 去除空格并转换为整数(num, num) // 生成(key, value)键值对key用于排序value保留原始数值}//第三步对RDD进行重新分区变成一个分区保证全局排序有序val partitionedRDD pairRDD.partitionBy(new HashPartitioner(1))//第四步执行sortByKey操作对RDD中所有元素都按照key的升序排序。val sortedRDD partitionedRDD.sortByKey(true) // true表示升序排序//第五步执行keys操作将键值对RDD中所有元素的key返回形成一个新的RDD。val sortedKeysRDD sortedRDD.keys//第六步执行map操作取出RDD中每个元素生成一个(key, value)键值对其中key是排序位次value是原整数。// 使用zipWithIndex生成位次索引从0开始需1val rankRDD sortedKeysRDD.zipWithIndex().map { case (num, index) (index 1, num) // 位次从1开始所以索引1作为key}//第七步执行collect操作以数组的形式返回RDD中所有元素。val resultArray rankRDD.collect()//第八步执行foreach操作依次遍历数组中每个元素按格式输出key valueprintln() // 保留原有打印行resultArray.foreach { case (rank, num) println(s$rank $num)}/********** End **********/}}第3关求平均值任务描述本关任务编写Spark独立应用程序实现求平均值。相关知识为了完成本关任务你需要掌握RDD的创建RDD的转换操作RDD的行动操作。编程要求每个输入文件表示学生某门课程的成绩输入文件中每行内容由两个字段组成第一个字段是学生姓名第二个字段是学生的成绩编写Spark独立应用程序求出所有学生的平均成绩并输出到一个新文件中。下面是输入文件和输出文件的一个样例供参考。输入文件(AlgorithmScore)的样例如下 XiaoMing 92 XiaoHong 87 XiaoXin 82 XiaoLi 90输入文件(DataBaseScore)的样例如下 XiaoMing 95 XiaoHong 81 XiaoXin 89 XiaoLi 85输入文件(PythonScore)的样例如下 XiaoMing 82 XiaoHong 83 XiaoXin 94 XiaoLi 91输出文件的样例如下 XiaoMing 89.67 XiaoXin 88.33 XiaoHong 83.67 XiaoLi 88.67import org.apache.spark.SparkContextimport org.apache.spark.SparkConfobject AvgScore {def main(args: Array[String]): Unit {val conf new SparkConf().setAppName(AvgScore).setMaster(local)val sc new SparkContext(conf)//输入文件AlgorithmScore.txt、DataBaseScore.txt和PythonScore.txt已保存在本地文件系统/root/step3_files目录中val dataFile file:///root/step3_filesval data sc.textFile(dataFile)/********** Begin **********///第一步执行过滤操作把空行丢弃。val filteredData data.filter(line line.trim.nonEmpty)//第二步执行map操作拆分每行文本为字符串数组姓名、成绩val splitRDD filteredData.map(line line.split(\\s)) // 按任意空白符拆分//第三步构建(key, value)键值对姓名成绩整数val pairRDD splitRDD.map { arr val name arr(0).trim // 学生姓名val score arr(1).trim.toInt // 成绩转换为整数(name, score)}//第四步执行mapValues操作将value转换为(成绩, 课程数)的元组val scoreCountRDD pairRDD.mapValues(score (score, 1))//第五步执行reduceByKey操作计算每个学生的总分和课程总数val totalRDD scoreCountRDD.reduceByKey { case ((sum1, count1), (sum2, count2)) (sum1 sum2, count1 count2) // 总分相加课程数相加}//第六步执行mapValues操作计算平均成绩保留两位小数val avgRDD totalRDD.mapValues { case (totalScore, totalCount) val avg totalScore.toDouble / totalCount // 计算平均分BigDecimal(avg).setScale(2, BigDecimal.RoundingMode.HALF_UP).toDouble // 保留两位小数}//第七步执行collect操作以数组形式返回所有元素val resultArray avgRDD.collect()//第八步执行foreach操作按格式打印姓名和平均成绩保留两位小数println() // 保留原有打印行resultArray.foreach { case (name, avgScore) println(f$name $avgScore%.2f)}/********** End **********/}}三十九、企业spark案例 —— 出租车轨迹分析第1关数据清洗任务描述本关任务将出租车轨迹数据规整化清洗掉多余的字符串。相关知识为了完成本关任务你需要掌握1.如何使用 SparkSQL 读取 CSV 文件2.如何使用正则表达式清洗掉多余字符串。SparkSQL 读取 CSVval spark SparkSession.builder().appName(Step1).master(local).getOrCreate()spark.read.option(header, true).option(delimiter, CSV分隔符).csv(文件存储的位置)option 参数说明header 为 true 将 CSV 第一行数据作为头部信息换一句来说就是将 CSV 的第一行数据作为 SparkSQL 表的字段delimiter 分隔符例如CSV 文件默认以英文逗号进行字段分隔那么 delimiter 为英文逗号如果文件以分号进行字段分隔那么 delimiter 为分号SparkSQL 自定义UDF函数用户定义函数User-defined functionsUDFs是大多数SQL环境的关键特性用于扩展系统的内置功能。UDF允许开发人员通过抽象其低级语言实现来在更高级语言如SQL中启用新功能。 Apache Spark也不例外并且提供了用于将UDF与Spark SQL工作流集成的各种选项。UDF对表中的单行进行转换以便为每行生成单个对应的输出值。例如大多数 SQL环境提供UPPER函数返回作为输入提供的字符串的大写版本。用户自定义函数可以在Spark SQL中定义和注册为UDF并且可以关联别名这个别名可以在后面的SQL查询中使用。编程要求在右侧编辑器补充代码将出租车轨迹数据规整化,清洗掉多余的字符串并使用 DataFrame.show() 打印输出。CSV文件内容如下清洗掉红框里面的 $ 、 字符由于这两字符出现的次数没有规律所以需要使用正则匹配。清洗后内容如下特别说明本案例的 CSV 文件是以 \t 进行字段分隔文件路径为 /root/data.csvimport org.apache.spark.sql.SparkSessionimport org.apache.spark.sql.functions._object Step1 {def main(args: Array[String]): Unit {val spark SparkSession.builder().appName(Step1).master(local).getOrCreate()/**********begin**********/// 1. 读取CSV文件分隔符为\t首行作为表头val df spark.read.option(header, true) // 将第一行作为列名.option(delimiter, \t) // 指定分隔符为制表符.csv(/root/data.csv)// 2. 定义UDF清洗$和字符正则匹配任意次数的$或并替换为空val cleanStrUdf udf((str: String) {if (str null || str.isEmpty) str // 处理null/空值else str.replaceAll([$], ) // 正则匹配$或任意次数并替换为空})// 3. 获取所有列名对每一列应用清洗UDF修正直接使用colName无需嵌套col()val columns df.columnsval cleanedDf columns.foldLeft(df) { (tempDf, colName) tempDf.withColumn(colName, cleanStrUdf(col(colName))) // 仅用col(colName)获取列}// 4. 展示清洗后的结果默认显示前20行cleanedDf.show()/**********end**********/spark.stop()}}第2关数据分析任务描述本关任务使用SparkSQL完成数据分析。相关知识为了完成本关任务你需要掌握如何使用SparkSQL进行数据分析FastJson 简述JSON 协议使用方便越来越流行JSON 的处理器有很多这里我介绍一下 FastJsonFastJson 是阿里的开源框架被不少企业使用是一个极其优秀的Json框架Github地址FastJson 。FastJson 优点FastJson 数度快无论序列化和反序列化都是当之无愧的fast功能强大支持普通JDK类包括任意Java Bean Class、Collection、Map、Date或enum零依赖没有依赖其它任何类库import com.alibaba.fastjson.JSONimport org.apache.spark.sql.SparkSessionimport org.apache.spark.sql.functions._import java.text.SimpleDateFormatimport java.util.{Date, Locale}object Step2 {def main(args: Array[String]): Unit {val spark SparkSession.builder().appName(Step1).master(local).getOrCreate()spark.sparkContext.setLogLevel(error)/**********begin**********/// 1. 读取CSV文件制表符分隔首行作为表头val df spark.read.option(header, true).option(delimiter, \t).csv(/root/data2.csv)// 1. 将时间戳转换成时间格式yyyy-MM-dd补零 // 定义时间戳转日期的UDF修正格式为yyyy-MM-dd确保补零val timestampToDateUdf udf((ts: String) {if (ts null || ts.isEmpty) nullelse {// 格式改为yyyy-MM-dd强制月份/日期为两位补零val sdf new SimpleDateFormat(yyyy-MM-dd, Locale.ENGLISH)sdf.format(new Date(ts.toLong * 1000)) // 秒级时间戳转毫秒}})val dfWithTime df.withColumn(TIME, timestampToDateUdf(col(TIMESTAMP)))// 2. 按预期分步展示结果匹配题目输出顺序 // 第一步展示仅含TIME的基础数据dfWithTime.select(TRIP_ID, CALL_TYPE, ORIGIN_CALL, TAXI_ID, ORIGIN_STAND,POLYLINE, TIME).show()// 2. 将POLYLINE字段分离出startLocation,endLocation 两个字段 // 解析POLYLINE为JSON数组提取第一个/最后一个元素val parsePolylineUdf udf((polyline: String) {if (polyline null || polyline.isEmpty || polyline []) {(null, null, 0) // (起点, 终点, 点数)} else {try {val jsonArray JSON.parseArray(polyline)val pointCount jsonArray.size()val start if (pointCount 0) jsonArray.getString(0) else nullval end if (pointCount 0) jsonArray.getString(pointCount - 1) else null(start, end, pointCount)} catch {case _: Exception (null, null, 0) // 异常处理}}})// 应用UDF并拆分出startLocation、endLocation、pointCountval dfWithLocation dfWithTime.withColumn(polyline_info, parsePolylineUdf(col(POLYLINE))).withColumn(startLocation, col(polyline_info._1)).withColumn(endLocation, col(polyline_info._2)).withColumn(pointCount, col(polyline_info._3).cast(int)).drop(polyline_info) // 删除临时列// 第二步展示含start/endLocation的数据dfWithLocation.select(TRIP_ID, CALL_TYPE, ORIGIN_CALL, TAXI_ID, ORIGIN_STAND,POLYLINE, TIME, startLocation, endLocation).show()// 3. 计算时长行程的总行程时间定义为点数-1×15秒 val dfWithTimeLen dfWithLocation.withColumn(timeLen, (col(pointCount) - 1) * 15)// 处理点数为0/1的情况时长为0.withColumn(timeLen, when(col(timeLen) 0, 0).otherwise(col(timeLen)))// 第三步展示含timeLen的完整数据dfWithTimeLen.select(TRIP_ID, CALL_TYPE, ORIGIN_CALL, TAXI_ID, ORIGIN_STAND,POLYLINE, TIME, startLocation, endLocation, timeLen).show()// 4. 统计每天各种呼叫类型的数量并以CALL_TYPE,TIME升序排序 val statDF dfWithTimeLen.groupBy(CALL_TYPE, TIME).count().withColumnRenamed(count, num) // 重命名为num.orderBy(col(CALL_TYPE).asc, col(TIME).asc)// 第四步展示统计结果statDF.show()/**********end**********/spark.stop()}}四十、RDD、DataSet 与 DataFrame 的转换Scala任务描述本关任务完成 RDD、DataSet 与 DataFrame 之间的相互转换。相关知识为了完成本关任务你需要掌握熟悉基础的 Scala 代码编写创建 RDD、DataSet、DataFrame 数据集。编程要求使用 Scala 编写工程代码根据所给 RDD先转换为 DataFrame 格式然后再将其转换为 DataSet 数据集格式并输出结果。任务说明 打开右侧代码文件窗口在 Begin 至 End 区域补充代码完善程序根据所给 RDD先转换为 DataFrame 格式然后再将其转换为 DataSet 数据集格式并输出结果。import org.apache.spark.rdd.RDDimport org.apache.spark.sql.{DataFrame, Dataset, SparkSession}object First_Question {case class Employee(id:Int,e_name:String,e_part:String,salary:Int)def main(args: Array[String]): Unit {val spark: SparkSession SparkSession.builder().appName(First_Question).master(local[*]).getOrCreate()val rdd: RDD[(Int, String, String, Int)] spark.sparkContext.parallelize(List((1001, 李晓, 运营部, 6000), (1002, 张花, 美术部, 6000), (1003, 李强, 研发部, 8000), (1004,田美, 营销部, 5000), (1005, 王菲, 后勤部, 4000)))/******************* Begin *******************/// 1. 导入隐式转换必须在SparkSession创建后导入import spark.implicits._// 2. 将RDD转换为DataFrame通过ScalaBean映射val df: DataFrame rdd.map(line {Employee(line._1, line._2, line._3, line._4)}).toDF()// 3. 将DataFrame转换为DataSet结合Employee Beanval ds: Dataset[Employee] df.as[Employee]// 4. 输出DataSet结果ds.show()/******************* End *******************/spark.stop()}}有任何问题都可以随时关注私信