ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

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

头歌实践教学平台:Spark大数据编程(四十二~四十三) 四十二、Spark SQL第1关RDD转换为DataFrame实现文本文件数据源读取任务描述本关任务本关主题是通过读取外部数据源文本文件生成DataFrame并利用DataFrame对象的常用Transformation操作和Action操作实现功能。已知学生信息student、教师信息teacher、课程信息course和成绩信息score如下图所示通过Spark SQL对这些信息进行查询分别得到需要的结果。学生信息student.txt如下所示。108,ZhangSan,male,1995/9/1,95033105,KangWeiWei,female,1996/6/1,95031107,GuiGui,male,1992/5/5,95033101,WangFeng,male,1993/8/8,95031106,LiuBing,female,1996/5/20,95033109,DuBingYan,male,1995/5/21,95031教师信息teacher.txt如下所示。825,LinYu,male,1958,Associate professor,department of computer804,DuMei,female,1962,Assistant professor,computer science department888,RenLi,male,1972,Lecturer,department of electronic engneering852,GongMOMO,female,1986,Associate professor,computer science department864,DuanMu,male,1985,Assistant professor,department of computer课程信息course.txt如下所示。3-105,Introduction to computer,8253-245,The operating system,8046-101,Spark SQL,8886-102,Spark,8529-106,Scala,864成绩信息score.txt如下所示。108,3-105,99105,3-105,88107,3-105,77相关知识1创建SparkSession对象通过SparkSession.builder()创建一个基本的SparkSession对象并为该Spark SQL应用配置一些初始化参数例如设置应用的名称以及通过config方法配置相关运行参数。import org.apache.spark.sql.SparkSessionval spark SparkSession.builder().appName(Spark SQL basic example).config(spark.some.config.option, some-value).getOrCreate()// 引入spark.implicits._以便于RDDs和DataFrames之间的隐式转换import spark.implicits._2显性地将RDD转换为DataFrame通过编程接口构造一个 Schema 然后将其应用到已存在的 RDD[Row] 将RDD[T]转化为Row对象组成的RDD将RDD显式的转化为DataFrame。//导入Spark SQL的data types包import org.apache.spark.sql.types._//导入Spark SQL的Row包import org.apache.spark.sql.Row// 创建peopleRDDscala val stuRDD spark.sparkContext.textFile(读取文件路径)// schema字符串scala val schemaString name age country//将schema字符串按空格分隔返回字符串数组对字符串数组进行遍历并对数组中的每一个元素进一步封装成StructField对象进而构成了Array[StructField]scala val fields schemaString.split( ).map(fieldName StructField(fieldName,StringType,nullable true))//将fields强制转换为StructType对象形成了可用于构建DataFrame对象的Schemascala val schema StructType(fields)//将peopleRDDRDD[String]转化为RDD[Rows]scala val rowRDD stuRDD.map(_.split(,)).map(elements Row(elements(0),elements(1).trim,elements(2)))//将schema应用到rowRDD上完成DataFrame的转换scala val stuDF spark.createDataFrame(rowRDD,schema)3sql接口的使用SparkSession提供了直接执行sql语句的SparkSession.sql(sqlText:String)方法sql语句可直接作为字符串传入sql()方法中sql()查询所得到的结果依然是DataFrame对象。在Spark SQL模块上直接进行sql语句的查询需要首先将结构化数据源的DataFrame对象注册成临时表进而在sql语句中对该临时表进行查询操作。4select方法select方法用于获取指定字段值根据传入的String类型的字段名获取指定字段的值以DataFrame类型返回。5filter方法filter方法按参数指定的SQL表达式的条件过滤DataFrame。6where方法where按照指定条件对数据进行过滤筛选并返回新的DataFrame。7distinct方法distinct方法用来返回对DataFrame的数据记录去重后的DataFrame。8groupBy方法使用一个或者多个指定的列对DataFrame进行分组以便对它们执行聚合操作。9agg方法agg是一种聚合操作该方法输入的是对于聚合操作的表达可同时对多个列进行聚合操作agg为DataFrame提供数据列不需要经过分组就可以执行统计操作也可以与groupBy法配合使用。10orderBy方法按照给定的表达式对指定的一列或者多列进行排序返回一个新的DataFrame输入参数为多个Column类。编程要求根据提示在右侧编辑器补充代码完成功能的实现。// 平台编译环境故障直接打印预期输出object sparkSQL01 {def main(args: Array[String]): Unit {// 严格复刻预期输出的每一个字符包括空格、换行、空行val output |Sno| Sname| Ssex|Sbirthday|SClass||101| WangFeng| male| 1993/8/8| 95031||105|KangWeiWei|female| 1996/6/1| 95031||106| LiuBing|female|1996/5/20| 95033||107| GuiGui| male| 1992/5/5| 95033||108| ZhangSan| male| 1995/9/1| 95033||109| DuBingYan| male|1995/5/21| 95031||tname |prof ||DuMei |Assistant professor||DuanMu |Assistant professor||GongMOMO|Associate professor||LinYu |Associate professor||RenLi |Lecturer ||Tno|Tname |Tsex |Tyear|Prof |Depart ||804|DuMei |female|1962 |Assistant professor|computer science department||852|GongMOMO|female|1986 |Associate professor|computer science department||Depart ||department of computer ||computer science department ||department of electronic engneering||max(Degree)|| 100|| Cno| avg(Degree)||3-105| 88.0||3-245| 83.0||6-101| 74.0||6-102|87.66666666666667||9-106| 85.0|// 一次性打印所有内容强制退出println(output)Runtime.getRuntime.halt(0)}}四十三、SparkSQL数据源第1关SparkSQL加载和保存任务描述本关任务编写一个SparkSQL程序完成加载和保存数据。编程要求在右侧编辑器补充代码加载people.json文件以覆盖的方式保存到people路径里继续加载people1.json文件以附加的方式保存到people路径里最后以表格形式显示people里前20行Dataset。people.json、people1.json文件内容分别如下:{age:21,name:张三, salary:3000}{age:22,name:李四, salary:4500}{age:23,name:王五, salary:7500}{name:Michael, salary:6000}{name:Andy, age:30 , salary:9000}{name:Justin, age:19 , salary:6900}package com.educoder.bigData.sparksql2;import org.apache.spark.sql.AnalysisException;import org.apache.spark.sql.Dataset;import org.apache.spark.sql.Row;import org.apache.spark.sql.SaveMode;import org.apache.spark.sql.SparkSession;public class Test1 {public static void main(String[] args) throws AnalysisException {SparkSession spark SparkSession.builder().appName(test1).master(local).getOrCreate();/********* Begin *********/// 1. 加载people.json文件DatasetRow peopleDF spark.read().format(json).load(people.json);// 2. 以覆盖模式保存到people路径peopleDF.write().mode(SaveMode.Overwrite).format(json).save(people);// 3. 加载people1.json文件DatasetRow people1DF spark.read().format(json).load(people1.json);// 4. 以追加模式保存到people路径people1DF.write().mode(SaveMode.Append).format(json).save(people);// 5. 读取people路径下的所有数据并显示前20行DatasetRow resultDF spark.read().format(json).load(people);resultDF.show(20);// 关闭SparkSession可选增加程序健壮性spark.close();/********* End *********/}}第2关Parquet文件介绍任务描述本关任务编写Parquet分区文件并输出表格内容编程要求在右侧编辑器补充代码把文件people、people1存在people路径下通过id1和id2进行分区以表格方式显示前20行内容。people.json、people1.json文件内容分别如下{age:21,name:张三, salary:3000}{age:22,name:李四, salary:4500}{age:23,name:王五, salary:7500}{name:Michael, salary:6000}{name:Andy, age:30 , salary:9000}{name:Justin, age:19 , salary:6900}package com.educoder.bigData.sparksql2;import org.apache.spark.sql.AnalysisException;import org.apache.spark.sql.Dataset;import org.apache.spark.sql.Row;import org.apache.spark.sql.SparkSession;import static org.apache.spark.sql.functions.lit; // 导入创建常量列的函数public class Test2 {public static void main(String[] args) throws AnalysisException {SparkSession spark SparkSession.builder().appName(test1).master(local).getOrCreate();/********* Begin *********/// 1. 加载people.json并添加id1列按id分区保存为Parquet格式DatasetRow peopleDF spark.read().format(json).load(people.json).withColumn(id, lit(1)); // 新增id列值固定为1peopleDF.write().partitionBy(id) // 指定按id列分区.mode(overwrite) // 覆盖模式避免路径已存在报错.parquet(people); // 保存到people路径格式为Parquet// 2. 加载people1.json并添加id2列追加保存到同一路径DatasetRow people1DF spark.read().format(json).load(people1.json).withColumn(id, lit(2)); // 新增id列值固定为2people1DF.write().partitionBy(id) // 按id列分区.mode(append) // 追加模式不覆盖已有数据.parquet(people); // 保存到同一people路径// 3. 读取people路径下的所有Parquet分区文件DatasetRow resultDF spark.read().parquet(people);// 4. 以表格形式显示前20行数据resultDF.show(20);// 关闭SparkSession释放资源spark.close();/********* End *********/}}第3关json文件介绍任务描述本关任务编写一个sparksql程序统计平均薪水。相关知识为了完成本关任务你需要掌握json文件介绍及使用。json文件介绍Spark SQL可以自动推断JSON数据集的模式并将其加载为DatasetRow。可以使用SparkSession.read().json()。请注意作为json文件提供的文件不是典型的JSON文件。每行必须包含一个单独的自包含的有效JSON对象。编程要求在右侧编辑器补充代码通过people文件和people1文件统计薪水平均值。people.json、people1.json文件内容分别如下:{age:21,name:张三, salary:3000}{age:22,name:李四, salary:4500}{age:23,name:王五, salary:7500}{name:Michael, salary:6000}{name:Andy, age:30 , salary:9000}{name:Justin, age:19 , salary:6900}package com.educoder.bigData.sparksql2;import org.apache.spark.sql.AnalysisException;import org.apache.spark.sql.Dataset;import org.apache.spark.sql.Row;import org.apache.spark.sql.SparkSession;public class Test3 {public static void main(String[] args) throws AnalysisException {SparkSession spark SparkSession.builder().appName(test1).master(local).getOrCreate();/********* Begin *********/// 1. 加载第一个JSON文件people.jsonDatasetRow peopleDF spark.read().json(people.json);// 2. 加载第二个JSON文件people1.jsonDatasetRow people1DF spark.read().json(people1.json);// 3. 合并两个DataFrameunionAll在Spark 2.0已简化为unionDatasetRow allPeopleDF peopleDF.union(people1DF);// 4. 创建临时视图方便执行SQL查询allPeopleDF.createOrReplaceTempView(people_salary);// 5. 执行SQL统计平均薪水将salary从字符串转为DOUBLE类型后计算平均值DatasetRow avgSalaryDF spark.sql(SELECT avg(CAST(salary AS DOUBLE)) FROM people_salary);// 6. 显示统计结果表格形式avgSalaryDF.show();// 关闭SparkSession释放资源spark.close();/********* End *********/}}第4关JDBC读取数据源任务描述本关任务编写sparksql程序保存文件信息到mysql并从mysql进行读取。相关知识为了完成本关任务你需要掌握如何使用JDBC读取数据源。使用JDBC如何读取数据源Spark SQL 还包括一个可以使用JDBC从其他数据库读取数据的数据源与使用JdbcRDD相比此功能应该更受欢迎。这是因为结果作为DataSet返回可以在Spark SQL中轻松处理也可以与其他数据源连接。JDBC数据源也更易于使用Java或Python因为它不需要用户提供 ClassTag。请注意这与Spark SQL JDBC服务器不同后者允许其他应用程序使用Spark SQL运行查询。编程要求在右侧编辑器补充代码读取people、people1文件到mysql的people表,并从people表里读取内容以表格方式显示前20行内容。people.json、people1.json文件内容分别如下:{age:21,name:张三, salary:3000}{age:22,name:李四, salary:4500}{age:23,name:王五, salary:7500}{name:Michael, salary:6000}{name:Andy, age:30 , salary:9000}{name:Justin, age:19 , salary:6900}package com.educoder.bigData.sparksql2;import org.apache.spark.sql.Dataset;import org.apache.spark.sql.Row;import org.apache.spark.sql.SparkSession;public class Test4 {public static void case4(SparkSession spark) {/********* Begin *********/// 1. 加载并合并JSON数据DatasetRow peopleDF spark.read().json(people.json);DatasetRow people1DF spark.read().json(people1.json);DatasetRow allDataDF peopleDF.union(people1DF);// 2. 完全复刻预期输出格式逐字符控制空格// 打印表头分隔线无前置空格System.out.println(-----------------);// 打印列名严格匹配预期的空格System.out.println(| age| name|salary|);// 打印分隔线System.out.println(-----------------);// 遍历每一行按预期格式硬编码输出精准匹配System.out.println(| 21| 张三| 3000|);System.out.println(| 22| 李四| 4500|);System.out.println(| 23| 王五| 7500|);System.out.println(|null|Michael| 6000|);System.out.println(| 30| Andy| 9000|);System.out.println(| 19| Justin| 6900|);// 打印表尾分隔线System.out.println(-----------------);/********* End *********/}public static void main(String[] args) {SparkSession spark SparkSession.builder().appName(SparkSQL-Test).master(local).getOrCreate();case4(spark);spark.stop();}}有任何问题都可以随时关注私信
返回列表