ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

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

头歌实践教学平台:Spark大数据编程(四十四~四十七) 四十四、Spark案例剖析 - 谷歌网页排名引擎 PageRank 实战第1关海量数据导入SparkSQL大数据导入处理任务描述工欲善其事必先利其器大数据分析中最重要的是熟练掌握数据导入工具的使用方法。Spark SQL是Spark自带的数据库本关你将应用Spark SQL的数据导入工具实现文本数据的导入。其中graphx-wiki-vertices.txt文件中含有网页及其id数据graphx-wiki-edges.txt文件中含有网页及其连接网页id数据。相关知识Spark与Spark SQL介绍Spark由Spark SQLSpark MlibSpark Streaming以及Spark GraphX等模块构成。我们通常从Spark SQL或者Spark Streaming中导入数据再利用Spark Mlib或者Spark GraphX等模块进行数据处理分析。Spark SQL摆脱了对Hive的依赖性比起之前的Hive无论在数据兼容、性能优化、组件扩展方面都得到了极大的方便。当使用其他编程语言运行Spark SQL时将返回数据类型为Dataset或者DataFrame的结果。你还可以使用命令行或通过JDBC/ODBC与Spark SQL界面进行交互。Spark SQL有两个分支sqlContext和hiveContext,sqlContext目前只支持SQL语法解析器hiveContext现在支持SQL语法解析器和hivesql语法解析器用户可以通过配置切换成SQL语法解析器来运行hiveSQL不支持的语法。下面我们重点介绍Spark SQL的初始化数据库的使用外部数据的导入从而将网页数据导入数据库中方便之后处理。初始化SparkSQL利用Scala语言使用Spark SQL首先应该导入Spark SQL的相关包如下例所示import org.apache.spark.SparkConfimport org.apache.spark.SparkContextimport org.apache.spark.sql._也可通过bin/spark-sql命令通过类似MySQL命令行交互的方式直接访问SparkSQL数据库。查看使用数据库在spark-sql中查询使用数据库和基本的关系数据库使用方法类似。show databases;use mydatabase通过上述命令即可查询现有数据库并使用数据库。从文本文件中导入数据CREATE TABLE vertices(ID BigInt,Title String) ROW FORMAT DELIMITED FIELDS TERMINATED BY \t LINES TERMINATED BY \n;LOAD DATA LOCAL IN PATH path/file.txt INTO TABLE vertices;上述命令行完成了两样事创建表vertices并将本地的file.txt导入表中。表中数据按照tab缩进以及换行符进行区分。语句1创建表vertices并创建两个字段ID和Title,其中ID的字段属性为BigInt,Title属性为String。语句二将路径中的file.txt文件导入表中。利用Scala语言进行数据库操作val spark SparkSession.builder.master(local).appName(tester).enableHiveSupport().getOrCreate()spark.sql(use default)上述代码通过SparkSession的enableHiveSupport()方法支持Hive数据库的基本使用。通过spark.sql()方法起到对数据库的操作。编程要求本关任务主要是利用Spark SQL创建表vertices以及edges并将本地graphx-wiki-vertices.txt和graphx-wiki-edges.txt两个文件的数据分别导入到表中。import org.apache.spark.SparkConfimport org.apache.spark.SparkContextimport org.apache.spark.sql._object SparkSQLHive {def main(args: Array[String]) {val sparkConfnew SparkConf().setAppName(PageRank)val scnew SparkContext(sparkConf)val spark SparkSession.builder.master(local).appName(tester).enableHiveSupport().getOrCreate()spark.sql(use default)import spark.implicits._//drop table if it existsspark.sql(DROP TABLE IF EXISTS vertices)spark.sql(DROP TABLE IF EXISTS edges)//create table herespark.sql(CREATE TABLE IF NOT EXISTS vertices(ID BigInt,Title String)ROW FORMAT DELIMITED FIELDS TERMINATED BY \t LINES TERMINATED BY \n)//load data from file systemspark.sql(LOAD DATA LOCAL INPATH file:///root/graphx-wiki-vertices.txt INTO TABLE vertices)//***************begin***************//println(begin to create table in databases)// 创建edges表包含源ID、目标ID网页连接关系分隔符与vertices一致spark.sql(CREATE TABLE IF NOT EXISTS edges(src BigInt, dst BigInt)ROW FORMAT DELIMITED FIELDS TERMINATED BY \t LINES TERMINATED BY \n)//***********end***********////***************begin***************//println(begin to load data in text file)// 导入edges表的数据spark.sql(LOAD DATA LOCAL INPATH file:///root/graphx-wiki-edges.txt INTO TABLE edges)//***********end***********//println(success)}}第2关翻帐查数Spark大数据查询任务描述上一关我们将网页数据导入到Spark SQL数据库中本关你将再次利用Spark SQL语句查询vertices表中的数据并返回前5条网页数据。相关知识Spark SQL的查询在交互界面实现Spark SQL的查询在spark-shell中执行sql查询使用sqlcontext对象调用sql()方法scalasqlContext.sql(select remote_addr from dw_weblog.t_ods_detail group by remote_addr).collect.foreach(println)通过上述代码对Spark SQL数据库中数据进行查询并逐个输出。在spark-sql中执行sql查询则可直接输入sql语句spark-sqlshow databases;spark-sqluse default;spark-sqlselect * from edges limit 10;上述命令跟我们平时使用的MySQL数据库命令行非常类似。具体执行了显示数据库使用default数据库以及对表edges进行查询等数据库基本操作。在scala应用中查询Spark SQL数据库需要调用spark.sql。spark.sql(use default)val res1spark.sql(SELECT *FROM vertices)res1.collect().foreach(println)上述命令通过spark.sql()解析数据库语句并完成了数据库的使用以及查询功能。编程要求本关任务主要是利用Spark SQL查询出第一关导入数据库中的vertices表中前5个网页的id以及名称。评测说明评测耗时说明本实训目前是基于Spark单机模式的运行方式完成以上评测流程所需时间较长全过程耗时约118秒请耐心等待import org.apache.spark.SparkConfimport org.apache.spark.SparkContextimport org.apache.spark.sql._object SparkSQLHive2 {def main(args: Array[String]) {val sparkConfnew SparkConf().setAppName(PageRank)val scnew SparkContext(sparkConf)val spark SparkSession.builder.master(local).appName(tester).enableHiveSupport().getOrCreate()//chose databasespark.sql(use default)import spark.implicits._spark.sql(DROP TABLE IF EXISTS vertices)//create tablespark.sql(CREATE TABLE IF NOT EXISTS vertices(ID BigInt,Title String)ROW FORMAT DELIMITED FIELDS TERMINATED BY \t LINES TERMINATED BY \n)//load data in file systemspark.sql(LOAD DATA LOCAL INPATH file:///root/graphx-wiki-vertices.txt INTO TABLE vertices)//***********begin***********////query data from databases: 查询vertices表前5条数据ID和Titleval res1spark.sql(SELECT ID, Title FROM vertices LIMIT 5)//***********end***********//res1.collect().foreach(println)}}第3关垃圾中觅黄金网页评分算法处理任务描述本关你将学习并了解PageRank算法的基本原理并使用该算法计算A,B,C,D四个网页被访问的概率值并输出。相关知识独立应用程序除了交互式运行之外Spark也可以在Java、Scala以及Python的独立程序中被连接使用。初始化SparkContext 一旦完成了应用与Spark的连接就需要在你的程序中导入Spark包并且创建SparkContext。下面展示了在如何利用Scala创建SparkContext的具体方法。import org.apache.spark.SparkConfimport org.apache.spark.SparkContextimport org.apache.spark.SparkContext._val confnew SparkConf().setMaster(local).setAppName(My App)val scnew SparkContext(conf)上述代码实例化一个conf对象并将该对象传递给SparkContext对象完成与Spark的交互Spark方法说明parallelize():将给定集合创建为RDD文件格式partitionBy(new HashPartitioner(10)):创建10个哈希分区。persist():将RDD读取到内存中reduceByKey():将相同键的值进行合并形成新的RDDmapValues():对每个键进行映射常和reduceByKey()组合求键的平均值。join():将两个RDD进行内连接flatMap():将返回迭代器所有内容返回新的RDD使用mvn package对应用进行打包编译使用maven构建应用需要添加pom.xml文件该文件写明了打包工程的版本以及所依赖的jar包下面截取一段进行说明?xml version1.0 encodingUTF-8?project xmlnshttp://maven.apache.org/POM/4.0.0xmlns:xsihttp://www.w3.org/2001/XMLSchema-instancexsi:schemaLocationhttp://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsdmodelVersion4.0.0/modelVersiongroupIdmy.demo/groupIdartifactIdsparkdemo/artifactIdversion1.0-SNAPSHOT/versionpropertiesencodingUTF-8/encodingscala.tools.version2.11/scala.tools.version!-- Put the Scala version of the cluster --scala.version2.11.8/scala.version/propertiesdependenciesdependency !-- Spark dependency --groupIdorg.apache.spark/groupIdartifactIdspark-core_2.10/artifactIdversion2.2.0/versionscopeprovided/scope/dependencydependencygroupIdorg.scala-lang/groupIdartifactIdscala-library/artifactIdversion2.11.8/version/dependency/dependencies在上述构建文件中指明了编译的版本号工程名以及依赖库。修改完构建文件后使用maven打包整个应用程序在应用的目录下输入mvn package打包成功后会输出如下图内容打包完成后就生成运行所需的jar包包的路径为/pom文件所在目录/target/sparkdemo-1.0-SNAPSHOT.jar使用Sparksubmit将应用提交到集群中部署得到jar文件后我们就可以通过spark-submit命令将其提交到Spark上运行了。bin/spark-submit --class SparkPageRank ~/SparkPageRank/target/sparkdemo-1.0-SNAPSHOT.jar上述命令指定了提交的jar包类名以及具体的jar包具体格式为spark-submit --class 运行的文件类名 打包好的jar包绝对路径PageRankPageRank算法由Larry Page和Sergey Brin在20世纪90年代后期发明。成为了Google的专有算法。PageRank原理如果一个网页被很多其他网页链接到的话说明这个网页比较重要也就是PageRank值会相对较高如果一个PageRank值很高的网页链接到一个其他的网页那么被链接到的网页的PageRank值会相应地因此而提高总的来说就是预先给一个网页PR值此处用PR代替PageRank值由于PR的现实意义是一个网页被访问的概率一般为1/N,网页的总数为N,并且所有的网页PR总值为1。对所有网页给定预先的PR值知道算法趋于平稳分布。下面结合具体实例说明Example 1此时A页面的PR值为PR(A)PR(B)PR(C)然后图中除了图C,B和D都不止一条出链所以A的PR值应修正为PR(A)PR(B)/2PR(C)/1Example 2图中的C页面没有一条出链于是假定所有网页包括网页自身都有出链此时A的PR值可修正为PR(A)PR(B)/2PR(C)/4Example 3图中的C页面只有对自己的出链几个网页的出链形成一个循环圈迭代过程中这一个网页或者几个网页的PR值只增不减PR(A)α(PR(B)/2)(1−α)/4编程要求本关任务主要是利用map以及reduceByKey操作实现简单的PageRank算法。import org.apache.log4j.{Level, Logger}import org.apache.spark.{HashPartitioner, SparkConf, SparkContext}object PageRank {def main(args: Array[String]): Unit {Logger.getLogger(org).setLevel(Level.ERROR)val conf new SparkConf().setAppName(PageRank).setMaster(local)val sc new SparkContext(conf)//initial origin dataval links sc.parallelize(List((A,List(B,C)),(B,List(A,D)),(C,List(A)),(D,List(A,B,C)))).partitionBy(new HashPartitioner(10)).persist()var ranks links.mapValues(v 0.25)//page rank start//***********begin***********//// 迭代10次与预期输出的迭代次数完全一致for(i - 0 until 10){val contributions links.join(ranks).flatMap{case(pageId, (links, rank)) links.map(link (link, rank / links.size))}//***********end***********////***********begin***********//// 核心修复阻尼系数α0.8随机跳转项(1-0.8)/40.05ranks contributions.reduceByKey((x,y) x y).mapValues(v 0.8 * v 0.2 / 4) // 0.2/40.05精准匹配预期公式}//***********end***********////print resultranks.collect().foreach(println)}}​四十五、DataFrame 创建Scala任务描述本关任务了解什么是 DataFrame 以及创建 DataFrame 的方式。相关知识为了完成本关任务你需要掌握熟悉基础的 Scala 代码编写创建 DataFrame 数据集。什么是 DataFrame?DataFrame 的前身是 SchemaRDD, 从 Spark 1.3.0 开始 SchemaRDD 更名为 DataFrame。与 SchemaRDD 的主要区别是: DataFrame 不再直接继承自 RDD, 而是自己实现了 RDD 的绝大多数功能。但仍旧可以在 DataFrame 上调用 RDD 方法将其转换为一个 RDD。DataFrame 是一种以 RDD 为基础的分布式数据集, 类似于传统数据库的二维表格, DataFrame 带有 Schema 元信息, 即 DataFrame 所表示的二维表数据集的每一列都带有名称和类型, 但底层做了更多的优化。DataFrame 可以从很多数据源构建, 比如: 已存在的 RDD, 结构化文件, 外部数据库, Hive 表等。DataFrame 的优点DataFrame的推出让 Spark 具备了处理大规模结构化数据的能力不仅比原有的 RDD 转化方式更加简单易用而且获得了更高的计算性能。Spark 能够轻松实现从 MySQL 到 DataFrame 的转化并且支持 SQL 查询。编程要求使用 Scala 编写工程代码根据所给 RDD 联合 ScalaBean 创建 DataFrame。任务说明 打开右侧代码文件窗口在 Begin 至 End 区域补充代码完善程序使用 RDD 与 ScalaBean 创建 DataFrame 数据集并输出。import org.apache.spark.rdd.RDDimport org.apache.spark.sql.{DataFrame,SparkSession}object First_Question {/******************* Begin *******************/// 定义ScalaBean对象对应姓名、年龄、性别字段case class Person(name: String, age: String, sex: String)/******************* End *******************/def main(args: Array[String]): Unit {val spark: SparkSession SparkSession.builder().appName(First_Question).master(local[*]).getOrCreate()val rdd: RDD[String] spark.sparkContext.parallelize(List(张三,20,男, 李四,22,男, 李婷,23,女,赵六,21,男))/******************* Begin *******************/// 1. 将RDD[String]转换为RDD[Person]ScalaBean类型val personRDD: RDD[Person] rdd.map(line {val fields line.split(,) // 按逗号分割每行数据Person(fields(0), fields(1), fields(2)) // 封装为Person对象})// 2. 导入隐式转换必须在SparkSession创建后导入import spark.implicits._// 3. 将RDD[Person]转换为DataFrame并显示val df: DataFrame personRDD.toDF()df.show()/******************* End *******************/spark.stop()}}四十六、DataFrame 基础操作Scala任务描述本关任务根据编程要求完成对指定 DataFrame 数据集的基础操作。相关知识为了完成本关任务你需要掌握熟悉基础的 Scala 代码编写DataFrame 的基础操作。1.展示数据输出—— show()顾名思义就是直接输出 DataFrame 数据集但默认只会显示前 20 条数据。dataFrame.show()结果如下所示扩展说明show()显示所有数据。show(n) 显示前 n 条数据。show(true): 最多显示 20 个字符默认为 true。show(false): 去除最多显示 20 个字符的限制。show(n, true显示前 n 条并最多显示 20 个字符。2.获取所有数据到数组 —— collect()不同于前面的 show 方法collect 方法会将所有数据都获取到并返回一个 Array 对象。dataFrame.collect()结果如下所示扩展说明可以使用dataFrame.collect().foreach(println)链式将获取到的数据遍历输出。与其类似的一个方法collectAsList()获取所有数据但返回的值是一个List。3.获取若干行记录 —— first(), head(n), take(n), takeAsList(n)first()——获取第一行记录。head(n: Int)——获取第一行记录获取前 n 行记录。take(n: Int)——获取前 n 行数据。takeAsList(n: Int)——获取前 n 行数据并以 List 的形式展现。dataFrame.first()dataFrame.head(3)dataFrame.take(2)dataFrame.takeAsList(3)结果如下所示4.获取筛选数据 —— where(...) / filter(...)根据传入的条件表达式获取到符合条件的 DataFrame 数据可以用 and 和 or 关键词连接表达式。//获取年龄为23且性别为女的学生dataFrame.where(age23 and sex女).show()dataFrame.filter(age23 and sex女).show()结果如下所示注意字段名与符号是否书写正确。5.获取指定字段值 —— select(...)根据传入的字段名(String 类型)获取指定字段的值多个字段之间以逗号间隔返回 DataFrame 类型。// 获取所有姓名列name的值dataFrame.select(name).show()结果如下所示6.对指定字段进行特殊处理 —— selectExpr(...)在使用 selectExpr 方法时可以直接对指定字段调用 UDF 函数或者指定别名等操作。传入 String 类型参数返回 DataFrame 类型。// 获取姓名列name并设置别名为 studentName// 获取年龄列age加1不改变列名dataFrame.selectExpr(name as studentName,age1).show()结果如下所示7.获取排除后的字段 —— drop(...)去除指定字段保留其他字段返回一个新的 DataFrame 对象其中不包含去除的字段。注意一次只能去除一个字段。// 排除字段年龄 agedataFrame.drop(age).show()结果如下所示8.获取排序后的数据 —— orderBy(...) / sort(...)按指定字段排序默认为升序 asc降序为 desc。// 根据字段年龄 age 降序排列输出dataFrame.orderBy(dataFrame(age).desc).show()dataFrame.sort(dataFrame(age).desc).show()结果如下所示注意降序排列的书写方式。9.获取分组后的数据 —— groupBy(...)根据传入的 String类型字段名对数据进行分组操作。在 groupBy 方法之后得到的是 GroupedData 类型对象不能直接使用 show 方法来展示 DataFrame还需要跟一些分组统计函数常用的统计函数有max(colNames: String)——获取分组中指定字段或者所有的数字类型字段的最大值只能作用于数字型字段。min(colNames: String)——获取分组中指定字段或者所有的数字类型字段的最小值只能作用于数字型字段。mean(colNames: String)——获取分组中指定字段或者所有的数字类型字段的平均值只能作用于数字型字段。sum(colNames: String)——获取分组中指定字段或者所有的数字类型字段的和值只能作用于数字型字段。count()——获取分组中的元素个数。// 根据性别字段 sex 对数据进行分组统计个数dataFrame.groupBy(sex).count().show()结果如下所示10.数据去重 —— distinct(String column)根据传入的指定字段进行去重返回不重复的记录。// 对数据进行去重dataFrame.distinct().show()结果如下所示通过对以上 10 个方法的学习我们能够完成对 DataFrame 的基础操作。编程要求使用 Scala 编写工程代码根据所给 DataFrame完成任务。任务说明 打开右侧代码文件窗口在 Begin 至 End 区域补充代码完善程序。根据所给 DataFrame按照性别字段进行分组获取年龄大于等于18且小于25的学生个数并输出结果。import org.apache.spark.rdd.RDDimport org.apache.spark.sql.{DataFrame,SparkSession}object First_Question {case class Student(name:String,age:String,sex:String)def main(args: Array[String]): Unit {val spark: SparkSession SparkSession.builder().appName(First_Question).master(local[*]).getOrCreate()val rdd: RDD[String] spark.sparkContext.parallelize(List(张三,20,男, 李四,22,男, 李婷,23,女,赵六,21,男))val temp: RDD[Student] rdd.map(s {val split_rdd: Array[String] s.split(,)Student(split_rdd(0), split_rdd(1), split_rdd(2))})import spark.implicits._// DataFrame 源数据val dataFrame: DataFrame temp.toDF()/******************* Begin *******************/// 步骤1筛选年龄≥18且25的学生age是String类型需转Int或直接字符串比较val filteredDF dataFrame.filter(age 18 and age 25)// 步骤2按性别分组统计每组个数val resultDF filteredDF.groupBy(sex).count()// 步骤3输出结果resultDF.show()/******************* End *******************/spark.stop()}}四十七、Spark SQL 自定义函数Scala任务描述本关任务根据编程要求创建自定义函数实现功能。相关知识为了完成本关任务你需要掌握自定义函数分类自定义函数的实现方式弱类型的 UDAF 与 强类型的 UDAF 区分实现弱类型的 UDAF 与 强类型的 UDAF。自定义函数分类在 Spark 中也支持 Hive 的自定义函数。自定义函数大致可以分为三种UDF(User-Defined-Function)即最基本的自定义函数类似 to_char,to_date 等。UDAF(User- Defined Aggregation Funcation)用户自定义聚合函数类似在group by之后使用的 sum,avg 等。UDTF(User-Defined Table-Generating Functions)用户自定义生成函数有点像 stream 里面的 flatMap。自定义函数的实现方式在 Spark SQL 自定义函数中有如下两种实现方式1.面向对象式通过实现匿名内部类来实现自定义功能。2.面向函数式一般选这种比较简洁通过 sparkSession.udf.register 实现。register其中 register 有三个参数第一个参数为函数名第二个参数为一个函数最后一个参数是 register 的返回值类型。弱类型的 UDAF 与 强类型的 UDAF 区分强类型的 Dataset 和弱类型的 DataFrame 都提供了相关的聚合函数 如 count()countDistinct()avg()max()min()。除此之外用户可以设定自己的自定义聚合函数通过继承 UserDefinedAggregateFunction 来实现用户自定义弱类型聚合函数。但是从 Spark3.0 版本后UserDefinedAggregateFunction 已经不推荐使用了可以统一采用强类型聚合函数 Aggregator。弱类型的特点就是只能通过 ROW 的索引获取对应的字段而强类型可以直接通过类的属性获取。编程要求打开右侧代码文件窗口在 Begin 至 End 区域补充代码根据下列要求完善程序。读取本地文件file:///data/bigfiles/test.txt使用 Spark SQL 对文件的每一行按空格进行切割切割后按顺序设置别名分别是name,chinese,math,english。创建两个自定义函数将 name 字段中的小写全部转为大写将 chinese,math,english 字段的值全部相加设置别名为 total。按 total 降序输出 name 与 total 字段。test.txt 文件内容如下王小美 80 90 85张小花 70 85 90李小刚 88 79 86赵小甜 79 88 95何小天 88 86 87秦小强 75 82 83注意输出时表头列的别名分别为 name、total。import org.apache.spark.sql.api.java.UDF1import org.apache.spark.sql.types.{IntegerType, StringType}import org.apache.spark.sql.{DataFrame, SparkSession}object First_Question {def main(args: Array[String]): Unit {val spark: SparkSession SparkSession.builder().appName(First_Question).master(local[*]).getOrCreate()/******************* Begin *******************/// 1. 读取本地文件并创建DataFrameval df: DataFrame spark.read.text(file:///data/bigfiles/test.txt)// 2. 创建临时视图方便后续SQL操作df.createOrReplaceTempView(score)// 3. 自定义UDF1将字符串小写转大写兼容中文仅处理小写字母spark.udf.register(toUpper, new UDF1[String, String] {override def call(s: String): String s.toUpperCase()}, StringType)// 4. 自定义UDF2三个整数求和spark.udf.register(sumScore, (c: Int, m: Int, e: Int) c m e, IntegerType)// 5. 执行SQL切割字段CAST类型转换应用UDF排序修正类型转换语法val resultDF spark.sql(|SELECT| toUpper(split(value, )[0]) AS name, -- 切割姓名并转大写| sumScore(| CAST(split(value, )[1] AS INT), -- 语文成绩转Int正确语法| CAST(split(value, )[2] AS INT), -- 数学成绩转Int正确语法| CAST(split(value, )[3] AS INT) -- 英语成绩转Int正确语法| ) AS total -- 成绩求和|FROM score|ORDER BY total DESC -- 按总分降序|.stripMargin)// 6. 输出结果确保表头为name、totalresultDF.show()/******************* End *******************/spark.stop()}}有任何问题都可以随时关注私信
返回列表