# Bigdata_Training_Camp **Repository Path**: zhxuankun/bigdata_training_camp ## Basic Information - **Project Name**: Bigdata_Training_Camp - **Description**: No description available - **Primary Language**: Unknown - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2021-07-17 - **Last Updated**: 2022-03-30 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # Bigdata_Training_Camp ## MapReduce 编程作业 ### 环境信息 - Scala 2.13.6 - JDK 1.8 - sbt 1.5.5 程序实现详见 0711 模块。 编译方式(sbt-shell): ```sbtshell $ task1 / assembly ``` 生成 jar 包路径: `target\scala-2.13\zhengxuankun_0711-assembly-1.0.jar`。 ### 作业流程 1.登录 47.101.72.185 ```bash ssh student@47.101.72.185 ``` 2.进入程序所在目录 ```bash cd /home/student/zhengxuankun/0711 ``` 3.启动作业并查看输出结果 ```bash # 清除输出路径 hdfs dfs -rm -r /user/student/zhengxuankun/res # 提交作业,并指定 reducer 个数为 1 以及作业名称 yarn jar zhengxuankun_0711-assembly-1.0.jar -Dmapreduce.job.reduces=1 -Dmapreduce.job.name="zhengxuankun" /user/student/zhengxuankun/data /user/student/zhengxuankun/res # 拉取 hdfs 输出结果到本地并查看 hdfs dfs -getmerge /user/student/zhengxuankun/res ./res.log cat res.log ``` ## Hadoop RPC 作业 ### 环境信息 - Scala 2.13.6 - JDK 1.8 - sbt 1.5.5 编译方式(sbt-shell): ```sbtshell $ task2 / assembly ``` 生成 jar 包路径: `target\scala-2.13\zhengxuankun_0718-assembly-1.0.jar`。 程序实现详见 0718 模块。 ### 作业流程 1.登录 47.101.72.185 ```bash ssh student@47.101.72.185 ``` 2.进入程序所在目录 ```bash cd /home/student/zhengxuankun/0718 ``` 3.启动作业: - 启动 Server 端:`java -cp "zhengxuankun_0718-assembly-1.0.jar:/opt/cloudera/parcels/CDH/lib/hadoop/client/*" server.Server` ![image-20210725134507604](https://gitee.com/zhxuankun/Image/raw/master/ARTS_Tips/image-20210725134507604.png) - 启动 Client 端:`java -cp "zhengxuankun_0718-assembly-1.0.jar:/opt/cloudera/parcels/CDH/lib/hadoop/client/*" client.Client` ![image-20210725134546049](https://gitee.com/zhxuankun/Image/raw/master/ARTS_Tips/image-20210725134546049.png) ## HBase 作业 ### 环境信息 - Scala 2.13.6 - JDK 1.8 - sbt 1.5.5 编译方式(sbt-shell): ```shell $ task3 / assembly ``` 生成 jar 包路径: `target\scala-2.13\zhengxuankun_0725-assembly-1.0.jar`。 程序实现详见 0725 模块。 ### 作业流程 1.登录 47.101.72.185 ```bash ssh student@47.101.72.185 ``` 2.进入程序所在目录 ```bash cd /home/student/zhengxuankun/0725 ``` 3.启动作业:`java -cp "zhengxuankun_0725-assembly-1.0.jar:/opt/cloudera/parcels/CDH/jars/*" HBasePractice` ![image-20210803002007986](https://gitee.com/zhxuankun/Image/raw/master/ARTS_Tips/image-20210803002007986.png) ## Hive 作业 SQL 以及输出结果如下: ``` USE hive_sql_test1; SET hive.cli.print.header = true; SET hive.resultset.use.unique.column.names = false; /* 简单:展示电影ID为2116这部电影各年龄段的平均影评分 */ SELECT u.age AS age, avg(r.rate) AS avg_rating FROM t_user u JOIN t_rating r ON r.movieid = '2116' AND r.userid = u.userid GROUP BY u.age; /* 执行结果: Stage-Stage-2: Map: 1 Reduce: 1 Cumulative CPU: 11.17 sec HDFS Read: 24607121 HDFS Write: 294 HDFS EC Read: 0 SUCCESS Total MapReduce CPU Time Spent: 11 seconds 170 msec OK age avg_rating 1 3.2941176470588234 18 3.3580246913580245 25 3.436548223350254 35 3.2278481012658227 45 2.8275862068965516 50 3.32 56 3.5 Time taken: 47.991 seconds, Fetched: 7 row(s) */ /* 中等:找出男性评分最高且评分次数超过50次的10部电影,展示电影名,平均影评分和评分次数 */ WITH target_movie AS ( SELECT r.movieid, avg(r.rate) AS avg_rate, count(rate) AS count_rate FROM t_user u JOIN t_rating r ON u.sex = 'M' AND u.userid = r.userid GROUP BY r.movieid HAVING count(rate) > 50 ORDER BY avg_rate DESC LIMIT 10 ) SELECT m.moviename, t.avg_rate, t.count_rate FROM t_movie m JOIN target_movie t ON m.movieid = t.movieid ORDER BY t.avg_rate DESC; /* 执行结果: Stage-Stage-2: Map: 1 Reduce: 1 Cumulative CPU: 14.81 sec HDFS Read: 24606494 HDFS Write: 68113 HDFS EC Read: 0 SUCCESS Stage-Stage-3: Map: 1 Reduce: 1 Cumulative CPU: 5.22 sec HDFS Read: 73694 HDFS Write: 401 HDFS EC Read: 0 SUCCESS Stage-Stage-5: Map: 1 Reduce: 1 Cumulative CPU: 4.3 sec HDFS Read: 10026 HDFS Write: 733 HDFS EC Read: 0 SUCCESS Total MapReduce CPU Time Spent: 24 seconds 330 msec OK moviename avg_rate count_rate Sanjuro (1962) 4.639344262295082 61 Godfather, The (1972) 4.583333333333333 1740 Seven Samurai (The Magnificent Seven) (Shichinin no samurai) (1954) 4.576628352490421 522 Shawshank Redemption, The (1994) 4.560625 1600 Raiders of the Lost Ark (1981) 4.520597322348094 1942 Usual Suspects, The (1995) 4.518248175182482 1370 Star Wars: Episode IV - A New Hope (1977) 4.495307167235495 2344 Schindler's List (1993) 4.49141503848431 1689 Paths of Glory (1957) 4.485148514851486 202 Wrong Trousers, The (1993) 4.478260869565218 644 Time taken: 167.261 seconds, Fetched: 10 row(s) */ /* 困难:找出影评次数最多的女士所给出最高分的10部电影的平均影评分,展示电影名和平均影评分(可使用多行SQL) */ -- 找出影评次数最多的女士 userid WITH target_user AS ( SELECT u.userid, count(u.userid) AS rate_cnt FROM t_user u RIGHT JOIN t_rating r ON u.userid = r.userid GROUP BY u.userid ORDER BY rate_cnt DESC LIMIT 1 ), highest_rating_movie AS ( SELECT r.movieid, r.rate FROM target_user u JOIN t_rating r ON u.userid = r.userid ORDER BY r.rate DESC LIMIT 10 ) -- 根据上一步获得的 10 部电影的 movie ,关联求得电影名以及平均分 SELECT m.moviename, t.rate_avg FROM t_movie m JOIN ( SELECT r.movieid, avg(r.rate) AS rate_avg FROM t_rating r JOIN highest_rating_movie h ON r.movieid = h.movieid GROUP BY r.movieid ) t ON m.movieid = t.movieid;/* 执行结果: Stage-Stage-2: Map: 1 Reduce: 1 Cumulative CPU: 12.19 sec HDFS Read: 24604585 HDFS Write: 130154 HDFS EC Read: 0 SUCCESS Stage-Stage-3: Map: 1 Reduce: 1 Cumulative CPU: 4.79 sec HDFS Read: 135072 HDFS Write: 116 HDFS EC Read: 0 SUCCESS Stage-Stage-18: Map: 1 Cumulative CPU: 6.42 sec HDFS Read: 24599762 HDFS Write: 64675 HDFS EC Read: 0 SUCCESS Stage-Stage-5: Map: 1 Reduce: 1 Cumulative CPU: 4.3 sec HDFS Read: 69595 HDFS Write: 296 HDFS EC Read: 0 SUCCESS Stage-Stage-14: Map: 1 Cumulative CPU: 8.74 sec HDFS Read: 24600176 HDFS Write: 455 HDFS EC Read: 0 SUCCESS Stage-Stage-7: Map: 1 Reduce: 1 Cumulative CPU: 2.8 sec HDFS Read: 5955 HDFS Write: 376 HDFS EC Read: 0 SUCCESS Stage-Stage-13: Map: 1 Cumulative CPU: 2.37 sec HDFS Read: 6005 HDFS Write: 657 HDFS EC Read: 0 SUCCESS Total MapReduce CPU Time Spent: 41 seconds 610 msec OK moviename rate_avg Dr. Strangelove or: How I Learned to Stop Worrying and Love the Bomb (1963) 4.4498902706656915 Raise the Red Lantern (1991) 4.1732851985559565 Great Dictator, The (1940) 4.032407407407407 Fantasia (1940) 3.904559915164369 High Noon (1952) 4.1786600496277915 Big Sleep, The (1946) 4.312384473197782 Bambi (1942) 3.738539898132428 Big Chill, The (1983) 3.8747016706443915 Bull Durham (1988) 3.837442922374429 Dog Day Afternoon (1975) 3.969934640522876 Time taken: 500.101 seconds, Fetched: 10 row(s) */ ``` ## Spark 作业 ### 环境信息 - Scala 2.12.14 - JDK 1.8 - sbt 1.5.5 - spark-3.1.2-bin-hadoop3.2 编译方式(sbt-shell): ```shell $ task5 / assembly ``` 程序实现详见 0815 模块。 ### 作业流程 **作业一**: 思路:在项目 resources 目录下按照作业文档创建内容相同的 0.txt、1.txt、2.txt 三个文件。 main class:task1.InvertedIndex 在 Idea 执行后结果如下: ``` 21/08/24 22:16:39 INFO scheduler.DAGScheduler: Job 0 finished: collect at InvertedIndex.scala:25, took 1.189826 s (a,ArrayBuffer((2,1), (3,1))) (banana,ArrayBuffer((2,1), (3,1))) (is,ArrayBuffer((0,2), (1,1), (2,1), (3,1))) (it,ArrayBuffer((0,2), (1,1), (2,1), (3,1))) (what,ArrayBuffer((0,1), (1,1))) 21/08/24 22:16:39 INFO spark.SparkContext: Invoking stop() from shutdown hook 21/08/24 22:16:39 INFO server.AbstractConnector: Stopped Spark@67db0897{HTTP/1.1, (http/1.1)}{0.0.0.0:4041} ``` **作业二**: 思路:由于环境问题只能跑 local 模式,暂时以 local[xxx] 的思路控制并发度 main class: task2.SparkDistcp ```shell spark-submit --class task2.SparkDistcp zhengxuankun_0815-assembly-1.0.jar "-i" "-m 3" "/usr/local/src/spark-3.1.1-bin-hadoop2.7/data/mllib" "/usr/local/src/spark-3.1.1-bin-hadoop2.7/data/mllib_target" ``` source 与 target 文件如下: ![image-20210829171049845](https://gitee.com/zhxuankun/Image/raw/master/ARTS_Tips/image-20210829171049845.png) 程序设置 RDD 分区数为 9,提交作业时指定并发度为 3,查看 WebUI 的 task 并发度如下: ![image-20210824221304635](https://gitee.com/zhxuankun/Image/raw/master/ARTS_Tips/image-20210824221304635.png) 执行后 target 文件如下: ![image-20210829171427952](https://gitee.com/zhxuankun/Image/raw/master/ARTS_Tips/image-20210829171427952.png) ## SparkSQL 作业 ### 环境信息 - Scala 2.12.14 - JDK 11 - sbt 1.5.5 - spark-3.1.2-bin-hadoop3.2 编译方式(sbt-shell): ```shell $ task6 / assembly ``` 程序实现详见 0829 模块。 ### 作业一 思路:按照如下步骤添加自定义语法 1. 在 SqlBase.g4 的 statement 下定义命令 `SHOW VERSION` 2. 补充未定义过**Keyword**(SPARK-KEYWORD-LIST-START),使得 ANTLR 认识该关键字 3. **nonRederved**(DEFAULT-NON-RESERVED-START)、**ansiNonReserved**(非标准sql)也定义该关键字 4. 执行 antlr4(maven 插件),生成新代码 5. 去到 SparkSqlParser,重写方法 `vistiShowVersion` ```scala override def visitShowVersion(ctx: ShowVersionContext): ShowVersionCommand = withOrigin(ctx) { ShowVersionCommand(ctx) } ``` 6. 到 `org.apache.spark.sql.execution.command` 下继承 `RunnableCommand` 接口并打印所需信息 ```scala import org.apache.spark.internal.Logging import org.apache.spark.sql.{Row, SparkSession} import org.apache.spark.sql.catalyst.parser.SqlBaseParser.ShowVersionContext case class ShowVersionCommand(ctx: ShowVersionContext) extends RunnableCommand with Logging { override def run(sparkSession: SparkSession): Seq[Row] = { logWarning(s"java.version: ${System.getProperty("java.version")}") logWarning(s"SPARK_VERSION: ${sparkSession.version}") Seq.empty[Row] } } ``` 7. 将 `scalastyle-config.xml` 的 `maxLineLength` 从 100 修改为 150,否则会报 `[error] java.lang.RuntimeException: Failing because of negative scalastyle result` 导致编译失败 8. 执行编译命令:`./build/sbt -Phive -Phive-thriftserver -DskipTests package` 编译完成后,启动 spark-sql 并执行自定义命令: ``` SHOW VERSION; ``` 分别打印出 Java 版本及 Spark 版本: ![image-20210907200353557](https://gitee.com/zhxuankun/Image/raw/master/ARTS_Tips/image-20210907200353557.png) ### 作业二 思路:查看规则的注释及其测试用例,反推出能够命中要求的 sql。 1. 建立一个分区 Parquet 表 student,有 id、name、class 字段,按照年级 grade 分区,并插入数据: ``` set hive.exec.dynamic.partition.mode=nonstrict; create table if not exists student(id int, name string, class int) using parquet partitioned by (grade int); insert into table student values(1, '张三', 1, 2020), (2, '李四', 1, 2020), (3, '王五', 2, 2020), (4, '赵六', 1, 2021); ``` 2. 开启日志提示 ``` set spark.sql.planChangeLog.level=WARN; ``` 3. 对于要求一,有如下 sql 及命中规则 ``` select t.id as newid from (select id, name from student where 1=1 except select id, name from student where name = '张三') t; -- 逻辑计划规则,针对 Expect SELECT 进行转换,规则路径:ReplaceExceptWithFilter -> CombineFilters === Applying Rule org.apache.spark.sql.catalyst.optimizer.ReplaceExceptWithFilter === Project [id#0 AS newid#38] Project [id#0 AS newid#38] !+- Except false +- Distinct ! :- Project [id#0, name#1] +- Filter NOT coalesce((name#1 = 张三), false) ! : +- Filter (1 = 1) +- Project [id#0, name#1] ! : +- Relation[id#0,name#1,class#2,grade#3] parquet +- Filter (1 = 1) ! +- Project [id#39, name#40] +- Relation[id#0,name#1,class#2,grade#3] parquet ! +- Filter (name#40 = 张三) ! +- Relation[id#39,name#40,class#41,grade#42] parquet -- 逻辑计划规则,对相同列有多个 Project 时进行合并 === Applying Rule org.apache.spark.sql.catalyst.optimizer.CollapseProject === !Project [id#0 AS newid#38] Aggregate [id#0, name#1], [id#0 AS newid#38] !+- Aggregate [id#0, name#1], [id#0] +- Project [id#0, name#1] ! +- Project [id#0, name#1] +- Filter ((1 = 1) AND NOT coalesce((name#1 = 张三), false)) ! +- Filter ((1 = 1) AND NOT coalesce((name#1 = 张三), false)) +- Relation[id#0,name#1,class#2,grade#3] parquet ! +- Relation[id#0,name#1,class#2,grade#3] parquet -- 逻辑计划规则,消除无意义的过滤规则 === Applying Rule org.apache.spark.sql.catalyst.optimizer.BooleanSimplification === Aggregate [id#0, name#1], [id#0 AS newid#38] Aggregate [id#0, name#1], [id#0 AS newid#38] +- Project [id#0, name#1] +- Project [id#0, name#1] ! +- Filter (true AND NOT coalesce((name#1 = 张三), false)) +- Filter NOT coalesce((name#1 = 张三), false) +- Relation[id#0,name#1,class#2,grade#3] parquet +- Relation[id#0,name#1,class#2,grade#3] parquet ``` 4. 对于要求二,有如下 sql 及命中规则: ``` select t.id, t.name, 1 as x from (select id, name, grade from student where 1=1 except select id, name, grade from student where class=2) t where t.name='张三' order by x; -- 逻辑计划规则,针对 Expect SELECT 进行转换,与 ReplaceExceptWithFilter 类似 === Applying Rule org.apache.spark.sql.catalyst.optimizer.ReplaceExceptWithAntiJoin === Sort [x#115 ASC NULLS FIRST], true Sort [x#115 ASC NULLS FIRST], true +- Project [id#0, name#1, 1 AS x#115] +- Project [id#0, name#1, 1 AS x#115] +- Filter (name#1 = 张三) +- Filter (name#1 = 张三) ! +- Except false +- Distinct ! :- Project [id#0, name#1, grade#3] +- Join LeftAnti, (((id#0 <=> id#116) AND (name#1 <=> name#117)) AND (grade#3 <=> grade#119)) ! : +- Filter (1 = 1) :- Project [id#0, name#1, grade#3] ! : +- Relation[id#0,name#1,class#2,grade#3] parquet : +- Filter (1 = 1) ! +- Project [id#116, name#117, grade#119] : +- Relation[id#0,name#1,class#2,grade#3] parquet ! +- Filter (class#118 = 2) +- Project [id#116, name#117, grade#119] ! +- Relation[id#116,name#117,class#118,grade#119] parquet +- Filter (class#118 = 2) ! +- Relation[id#116,name#117,class#118,grade#119] parquet -- 逻辑计划规则,将 Dinstinct 替换为 Aggregate 操作 === Applying Rule org.apache.spark.sql.catalyst.optimizer.ReplaceDistinctWithAggregate === Sort [x#115 ASC NULLS FIRST], true Sort [x#115 ASC NULLS FIRST], true +- Project [id#0, name#1, 1 AS x#115] +- Project [id#0, name#1, 1 AS x#115] +- Filter (name#1 = 张三) +- Filter (name#1 = 张三) ! +- Distinct +- Aggregate [id#0, name#1, grade#3], [id#0, name#1, grade#3] +- Join LeftAnti, (((id#0 <=> id#116) AND (name#1 <=> name#117)) AND (grade#3 <=> grade#119)) +- Join LeftAnti, (((id#0 <=> id#116) AND (name#1 <=> name#117)) AND (grade#3 <=> grade#119)) :- Project [id#0, name#1, grade#3] :- Project [id#0, name#1, grade#3] : +- Filter (1 = 1) : +- Filter (1 = 1) : +- Relation[id#0,name#1,class#2,grade#3] parquet : +- Relation[id#0,name#1,class#2,grade#3] parquet +- Project [id#116, name#117, grade#119] +- Project [id#116, name#117, grade#119] +- Filter (class#118 = 2) +- Filter (class#118 = 2) +- Relation[id#116,name#117,class#118,grade#119] parquet +- Relation[id#116,name#117,class#118,grade#119] parquet -- 逻辑计划规则,谓词下推 === Applying Rule org.apache.spark.sql.catalyst.optimizer.PushDownPredicates === Sort [x#115 ASC NULLS FIRST], true Sort [x#115 ASC NULLS FIRST], true +- Project [id#0, name#1, 1 AS x#115] +- Project [id#0, name#1, 1 AS x#115] ! +- Filter (name#1 = 张三) +- Aggregate [id#0, name#1, grade#3], [id#0, name#1, grade#3] ! +- Aggregate [id#0, name#1, grade#3], [id#0, name#1, grade#3] +- Join LeftAnti, (((id#0 <=> id#116) AND (name#1 <=> name#117)) AND (grade#3 <=> grade#119)) ! +- Join LeftAnti, (((id#0 <=> id#116) AND (name#1 <=> name#117)) AND (grade#3 <=> grade#119)) :- Project [id#0, name#1, grade#3] ! :- Project [id#0, name#1, grade#3] : +- Filter ((1 = 1) AND (name#1 = 张三)) ! : +- Filter (1 = 1) : +- Relation[id#0,name#1,class#2,grade#3] parquet ! : +- Relation[id#0,name#1,class#2,grade#3] parquet +- Project [id#116, name#117, grade#119] ! +- Project [id#116, name#117, grade#119] +- Filter (class#118 = 2) ! +- Filter (class#118 = 2) +- Relation[id#116,name#117,class#118,grade#119] parquet ! +- Relation[id#116,name#117,class#118,grade#119] parquet -- 逻辑计划规则,尽可能的将属性名替换为可折叠表达式 === Applying Rule org.apache.spark.sql.catalyst.optimizer.FoldablePropagation === !Sort [x#115 ASC NULLS FIRST], true Sort [1 ASC NULLS FIRST], true +- Aggregate [id#0, name#1, grade#3], [id#0, name#1, 1 AS x#115] +- Aggregate [id#0, name#1, grade#3], [id#0, name#1, 1 AS x#115] +- Project [id#0, name#1, grade#3] +- Project [id#0, name#1, grade#3] +- Join LeftAnti, (((id#0 <=> id#116) AND (name#1 <=> name#117)) AND (grade#3 <=> grade#119)) +- Join LeftAnti, (((id#0 <=> id#116) AND (name#1 <=> name#117)) AND (grade#3 <=> grade#119)) :- Project [id#0, name#1, grade#3] :- Project [id#0, name#1, grade#3] : +- Filter ((1 = 1) AND (name#1 = 张三)) : +- Filter ((1 = 1) AND (name#1 = 张三)) : +- Relation[id#0,name#1,class#2,grade#3] parquet : +- Relation[id#0,name#1,class#2,grade#3] parquet +- Project [id#116, name#117, grade#119] +- Project [id#116, name#117, grade#119] +- Filter (class#118 = 2) +- Filter (class#118 = 2) +- Relation[id#116,name#117,class#118,grade#119] parquet +- Relation[id#116,name#117,class#118,grade#119] parquet -- 逻辑计划规则,常量折叠 === Applying Rule org.apache.spark.sql.catalyst.optimizer.ConstantFolding === Sort [1 ASC NULLS FIRST], true Sort [1 ASC NULLS FIRST], true +- Aggregate [id#0, name#1, grade#3], [id#0, name#1, 1 AS x#115] +- Aggregate [id#0, name#1, grade#3], [id#0, name#1, 1 AS x#115] +- Project [id#0, name#1, grade#3] +- Project [id#0, name#1, grade#3] +- Join LeftAnti, (((id#0 <=> id#116) AND (name#1 <=> name#117)) AND (grade#3 <=> grade#119)) +- Join LeftAnti, (((id#0 <=> id#116) AND (name#1 <=> name#117)) AND (grade#3 <=> grade#119)) :- Project [id#0, name#1, grade#3] :- Project [id#0, name#1, grade#3] ! : +- Filter ((1 = 1) AND (name#1 = 张三)) : +- Filter (true AND (name#1 = 张三)) : +- Relation[id#0,name#1,class#2,grade#3] parquet : +- Relation[id#0,name#1,class#2,grade#3] parquet +- Project [id#116, name#117, grade#119] +- Project [id#116, name#117, grade#119] +- Filter (class#118 = 2) +- Filter (class#118 = 2) +- Relation[id#116,name#117,class#118,grade#119] parquet +- Relation[id#116,name#117,class#118,grade#119] parquet ``` ### 作业三 思路:自定义规则中只打印出日志,不做任何转换。 将编译好的包拷贝到 Spark 目录下,指定 jar 包以及自定义规则类,启动 spark-sql: ```bash ./bin/spark-sql --jars zhengxuankun_0829-assembly-1.0.jar --conf spark.sql.extensions=task3.rules.MySparkSessionExtention ``` 执行如下语句即可看到规则打印的日志: ```bash set spark.sql.planChangeLog.level=WARN; select id, name from student where 1=1 ; 21/09/05 22:30:13 WARN MyRule: 优化规则, noop by zkx 21/09/05 22:30:13 WARN MyRule: === Applying Rule task3.rules.MyRule === ``` ## SparkSQL 优化作业 ### 环境信息 - Scala 2.12.14 - JDK 11 - sbt 1.5.5 - spark-3.1.2-bin-hadoop3.2 编译方式(sbt-shell): ```shell $ task7 / assembly ``` 程序实现详见 0908 模块。 ### 作业一 避免小文件的方式: 1. 在最终输出数据时,减少 task 的数目,如 repartition 2. 如果最终输出前有 Shuffle 阶段,可以使用 AQE 来动态调整 partition 数 3. 使用 hadoop har 工具 4. 如果 Spark 写入到 Hive orc 表,还可以使用 Hive 的 `alter table table_name partition(...) CONCATENATE ;` 合并分区下的小文件 5. 使用 hudi 提供的小文件合并功能 ### 作业二 在 SqlBase.g4 添加如下内容: ``` // statement 添加如下内容 | COMPACT TABLE tableIdentifier partitionSpec? (INTO INTEGER_VALUE FILES)? #compactTable // --ANSI-NON-RESERVED-START 添加如下内容 | FILES // --DEFAULT-NON-RESERVED-START 添加如下内容 | FILES // --SPARK-KEYWORD-LIST-START 添加如下内容 FILES: 'FILES'; ``` 执行 antlr 生成代码后,在如下文件中编写代码: ``` // org\apache\spark\sql\execution\SparkSqlParser.scala 添加如下内容 override def visitCompactTable(ctx: CompactTableContext): LogicalPlan = withOrigin(ctx) { val table = visitTableIdentifier(ctx.tableIdentifier()) // 文件个数 val filesNum = if (ctx.INTEGER_VALUE() != null) { Some(ctx.INTEGER_VALUE().getText) } else { None } CompactTableCommand(table, filesNum) } // org\apache\spark\sql\execution\command\tables.scala case class CompactTableCommand( table: TableIdentifier, filesNum: Option[String]) extends LeafRunnableCommand { // 临时表名称 val tempTable: TableIdentifier = TableIdentifier(s"temp_table_${System.nanoTime()}") override def output: Seq[Attribute] = Seq( AttributeReference("compat_stmt", StringType, nullable = false)() ) override def run(sparkSession: SparkSession): Seq[Row] = { sparkSession.catalog.setCurrentDatabase(table.database.getOrElse("default")) val tmpDF = sparkSession.table(table.identifier) val partitions = filesNum match { case Some(files) => files.toInt // 如果未指定文件数,则通过执行计划获取文件大小,并除以 128M 得到分区数 case None => (sparkSession.sessionState .executePlan(tmpDF.queryExecution.logical) .optimizedPlan .stats .sizeInBytes >> 27).toInt } // 写入临时表 tmpDF.repartition(partitions) .write .mode(SaveMode.Overwrite) .saveAsTable(tempTable.identifier) // 合并文件 sparkSession.table(tempTable.identifier) .write .mode(SaveMode.Overwrite) .saveAsTable(table.identifier) // 清除临时表 sparkSession.sql(s"DROP TABLE ${tempTable.identifier}") Seq(Row(s"Compacte Table $table finished.")) } } ``` 执行 `./build/sbt package -Phive -Phive-thriftserver -DskipTests` 进行编译。 准备好表数据如下: ![image-20210918145332770](https://gitee.com/zhxuankun/Image/raw/master/blog/image-20210918145332770.png) 拆分为 8 个文件: ```sql COMPACT TABLE t2 INTO 8 FILES; ``` ![image-20210918153840545](https://gitee.com/zhxuankun/Image/raw/master/blog/image-20210918153840545.png) 不指定文件数: ```SQL COMPACT TABLE t2; ``` ![image-20210920115349887](https://gitee.com/zhxuankun/Image/raw/master/ARTS_Tips/image-20210920115349887.png) ### 作业三 思路:在执行 insert 语句时,在 Scan 与 InsertInto 之间添加一个 Exchange 节点,合并为 1 个文件。在不修改 Spark 任何默认配置的前提下,发现 spark-sql 都是使用的 Hive 方式,因此规则也按照 Hive 节点来编写。 启动程序 `./bin/spark-sql --jars zhengxuankun_0908-assembly-1.0.jar --conf spark.sql.extensions=task3.rules.MySparkSessionExtention ` 并执行如下语句: ```sql insert into t3 select * from t2; ``` ![image-20210920131139648](https://gitee.com/zhxuankun/Image/raw/master/ARTS_Tips/image-20210920131139648.png) ![image-20210920131202992](https://gitee.com/zhxuankun/Image/raw/master/ARTS_Tips/image-20210920131202992.png) ![image-20210920131216753](https://gitee.com/zhxuankun/Image/raw/master/ARTS_Tips/image-20210920131216753.png) > TODO: 研究如何添加 Exchange 节后,再添加一个已有的 AQE 合并小文件节点。 ## Presto 作业 ### 环境信息 - Scala 2.13.6 - JDK 11 - sbt 1.5.5 编译方式(sbt-shell): ```shell $ task8 / assembly ``` 程序实现详见 0919 模块。 ### 作业一 1. 对于微博这种统计一些热门评论点赞数的,不需要结果太精确,可以使用 HyperLogLog 2. 电商大促统计某些页面的访问量 3. Postgres 支持以 extension 的方式使用 Hyperloglog 统计数据元信息 4. Redis 中使用 HyperLogLog,实现数据集合并时的去重 5. Presto、Kylin 使用 HyperLogLog 快速统计 Cardinality ### 作业二 使用本地 Docker 搭建 Presto 的方式执行作业,Docker 镜像地址:[ahanaio/prestodb-sandbox](https://hub.docker.com/r/ahanaio/prestodb-sandbox)。 1. 拉取镜像到本地并启动: ```bash docker run -p 8080:8080 --name presto ahanaio/prestodb-sandbox ``` 2. 另起一个 Session,连接上 Presto-cli: ```bash docker exec -it presto presto-cli ``` 3. 使用 tpcds.sf100 作为测试数据集,item 作为测试表。 4. 使用 HyperLogLog 方式统计 tpcds.sf100.item 中 i_item_sk 字段的 Cardinality: ```sql presto:sf100> select cardinality(approx_set(i_item_sk)) from tpcds.sf100.item; _col0 -------- 204126 (1 row) Query 20210929_094140_00045_n3i9g, FINISHED, 1 node Splits: 41 total, 41 done (100.00%) 0:07 [204K rows, 0B] [30.5K rows/s, 0B/s] ``` 5. 再使用 count(distinct) 语句精确统计 tpcds.sf100.item 中 i_item_sk 字段的 Cardinality: ```sql presto:sf100> select count(distinct(i_item_sk)) from item; _col0 -------- 204000 (1 row) WARNING: COUNT(DISTINCT xxx) can be a very expensive operation when the cardinality is high for xxx. In most scenarios, using approx_distinct instead would be enough Query 20210929_094400_00046_n3i9g, FINISHED, 1 node Splits: 57 total, 57 done (100.00%) 0:07 [204K rows, 0B] [29.5K rows/s, 0B/s] ``` 结论:可以看到,HyperLogLog 统计出来的 Cardinality 与实际相差不到千分之一。 参考链接:[HyperLogLog Functions](https://prestodb.io/docs/current/functions/hyperloglog.html#functions-hyperloglog--page-root) ### 作业三 思路:代码中执行与作业二同样的 SQL 语句,并打印 ResultSet 的内容。 执行 task3.presto.PrestoCli,可以看到打印出了如下结果: ``` HyperLogLog: _col0 204126 Count(Distinct): _col0 204000 ``` ## Flink 作业 SpendReport.java 的实现如下: ```java public static Table report(Table transactions) { Table table = transactions .window(Slide.over(lit(5).hours()) .every(lit(1).hours()) .on($("transaction_time")) .as("log_ts")) .groupBy($("log_ts"), $("account_id")) .select($("account_id"), $("log_ts").start(), $("amount").avg().as("amount")); table.printSchema(); return table; } ``` MySQL 结果: ![image-20211023235126957](https://gitee.com/zhxuankun/Image/raw/master/ARTS_Tips/image-20211023235126957.png) Grafana 展示: ![image-20211023235014450](https://gitee.com/zhxuankun/Image/raw/master/ARTS_Tips/image-20211023235014450.png) ## 毕业项目 ### 作业一:分析一条 TPCDS SQL(请基于 Spark 3.1.1 版本解答) 执行步骤如下: ```shell git clone https://github.com/maropu/spark-tpcds-datagen.git cd spark-tpcds-datagen axel https://archive.apache.org/dist/spark/spark-3.1.1/spark-3.1.1-bin-hadoop2.7.tgz tar -zxvf spark-3.1.1-bin-hadoop2.7.tgz mkdir -p tpcds-data-1g export SPARK_HOME=./spark-3.1.1-bin-hadoop2.7 ./bin/dsdgen --output-location tpcds-data-1g axel https://repo1.maven.org/maven2/org/apache/spark/spark-catalyst_2.12/3.1.1/spark-catalyst_2.12-3.1.1-tests.jar axel https://repo1.maven.org/maven2/org/apache/spark/spark-core_2.12/3.1.1/spark-core_2.12-3.1.1-tests.jar axel https://repo1.maven.org/maven2/org/apache/spark/spark-sql_2.12/3.1.1/spark-sql_2.12-3.1.1-tests.jar # 执行 q38 查询 ./spark-3.1.1-bin-hadoop2.7/bin/spark-submit --class org.apache.spark.sql.execution.benchmark.TPCDSQueryBenchmark --jars spark-core_2.12-3.1.1-tests.jar,spark-catalyst_2.12-3.1.1-tests.jar --conf spark.sql.planChangeLog.level=WARN spark-sql_2.12-3.1.1-tests.jar --data-location tpcds-data-1g --query-filter "q38" ``` #### 优化规则一:ColumnPruning 查看日志,遇到第一条优化规则如下: ``` 21/11/14 11:29:59 WARN PlanChangeLogger: === Applying Rule org.apache.spark.sql.catalyst.optimizer.ColumnPruning === Aggregate [count(1) AS count#1710L] Aggregate [count(1) AS count#1710L] !+- Relation[ib_income_band_sk#1697,ib_lower_bound#1698,ib_upper_bound#1699] parquet +- Project ! +- Relation[ib_income_band_sk#1697,ib_lower_bound#1698,ib_upper_bound#1699] parquet ``` ColumnPruning 属于 **RewriteSubquery** 规则组,从计划前后变化来看,ColumnPruning 在 Aggregate 节点与 Relation 节点之间插入了一个 Project 节点,对应源码为: ```scala object ColumnPruning extends Rule[LogicalPlan] { def apply(plan: LogicalPlan): LogicalPlan = removeProjectBeforeFilter(plan transform { // ...... // Prunes the unused columns from child of Aggregate/Expand/Generate/ScriptTransformation case a @ Aggregate(_, _, child) if !child.outputSet.subsetOf(a.references) => a.copy(child = prunedChild(child, a.references)) // ...... } /** Applies a projection only when the child is producing unnecessary attributes */ private def prunedChild(c: LogicalPlan, allReferences: AttributeSet) = if (!c.outputSet.subsetOf(allReferences)) { Project(c.output.filter(allReferences.contains), c) } else { c } ``` 在 Relation 与 Aggregate 之间插入 Project 节点,裁剪掉 Aggregate 不包含的列。 需要注意 apply 方法中的高阶函数:`removeProjectBeforeFilter` 。当遇到 `Project(_, f @ Filter(_, p2 @ Project(_, child)))` 的组合,即 Filter 节点的 child 也包含一个 Project节点,并且该 Project 节点的输出是 Project 节点自身 child 的子集,那么将该 Project 节点中的 child 外提,直接作为 Filter 的 child,消去 Project 节点: ```scala private def removeProjectBeforeFilter(plan: LogicalPlan): LogicalPlan = plan transformUp { case p1 @ Project(_, f @ Filter(_, p2 @ Project(_, child))) if p2.outputSet.subsetOf(child.outputSet) && // We only remove attribute-only project. p2.projectList.forall(_.isInstanceOf[AttributeReference]) => p1.copy(child = f.copy(child = child)) } ``` 这么做的原因主要是由于在 Filter 之前做 Project 并不是必须的,并且与谓词下推穿过 Project 相冲突。 #### 优化规则二:ReplaceIntersectWithSemiJoin 第二条优化规则如下: ``` === Applying Rule org.apache.spark.sql.catalyst.optimizer.ReplaceIntersectWithSemiJoin === OverwriteByExpression RelationV2[] noop-table, true, true OverwriteByExpression RelationV2[] noop-table, true, true +- GlobalLimit 100 +- GlobalLimit 100 +- LocalLimit 100 +- LocalLimit 100 +- Aggregate [count(1) AS count(1)#2008L] +- Aggregate [count(1) AS count(1)#2008L] ! +- Intersect false +- Distinct ! :- Intersect false +- Join LeftSemi, (((c_last_name#163 <=> c_last_name#1998) AND (c_first_name#162 <=> c_first_name#1997)) AND (d_date#331 <=> d_date#1963)) ! : :- Distinct :- Distinct ! : : +- Project [c_last_name#163, c_first_name#162, d_date#331] : +- Join LeftSemi, (((c_last_name#163 <=> c_last_name#1952) AND (c_first_name#162 <=> c_first_name#1951)) AND (d_date#331 <=> d_date#1917)) ! : : +- Filter (((ss_sold_date_sk#1154 = d_date_sk#329) AND (ss_customer_sk#1157 = c_customer_sk#154)) AND ((d_month_seq#332 >= 1200) AND (d_month_seq#332 <= (1200 + 11)))) : :- Distinct ! : : +- Join Inner : : +- Project [c_last_name#163, c_first_name#162, d_date#331] ! : : :- Join Inner : : +- Filter (((ss_sold_date_sk#1154 = d_date_sk#329) AND (ss_customer_sk#1157 = c_customer_sk#154)) AND ((d_month_seq#332 >= 1200) AND (d_month_seq#332 <= (1200 + 11)))) ! : : : :- Relation[ss_sold_date_sk#1154,ss_sold_time_sk#1155,ss_item_sk#1156,ss_customer_sk#1157,ss_cdemo_sk#1158,ss_hdemo_sk#1159,ss_addr_sk#1160,ss_store_sk#1161,ss_promo_sk#1162,ss_ticket_number#1163,ss_quantity#1164,ss_wholesale_cost#1165,ss_list_price#1166,ss_sales_price#1167,ss_ext_discount_amt#1168,ss_ext_sales_price#1169,ss_ext_wholesale_cost#1170,ss_ext_list_price#1171,ss_ext_tax#1172,ss_coupon_amt#1173,ss_net_paid#1174,ss_net_paid_inc_tax#1175,ss_net_profit#1176] parquet : : +- Join Inner ! : : : +- Relation[d_date_sk#329,d_date_id#330,d_date#331,d_month_seq#332,d_week_seq#333,d_quarter_seq#334,d_year#335,d_dow#336,d_moy#337,d_dom#338,d_qoy#339,d_fy_year#340,d_fy_quarter_seq#341,d_fy_week_seq#342,d_day_name#343,d_quarter_name#344,d_holiday#345,d_weekend#346,d_following_holiday#347,d_first_dom#348,d_last_dom#349,d_same_day_ly#350,d_same_day_lq#351,d_current_day#352,... 4 more fields] parquet : : :- Join Inner ! : : +- Relation[c_customer_sk#154,c_customer_id#155,c_current_cdemo_sk#156,c_current_hdemo_sk#157,c_current_addr_sk#158,c_first_shipto_date_sk#159,c_first_sales_date_sk#160,c_salutation#161,c_first_name#162,c_last_name#163,c_preferred_cust_flag#164,c_birth_day#165,c_birth_month#166,c_birth_year#167,c_birth_country#168,c_login#169,c_email_address#170,c_last_review_date#171] parquet : : : :- Relation[ss_sold_date_sk#1154,ss_sold_time_sk#1155,ss_item_sk#1156,ss_customer_sk#1157,ss_cdemo_sk#1158,ss_hdemo_sk#1159,ss_addr_sk#1160,ss_store_sk#1161,ss_promo_sk#1162,ss_ticket_number#1163,ss_quantity#1164,ss_wholesale_cost#1165,ss_list_price#1166,ss_sales_price#1167,ss_ext_discount_amt#1168,ss_ext_sales_price#1169,ss_ext_wholesale_cost#1170,ss_ext_list_price#1171,ss_ext_tax#1172,ss_coupon_amt#1173,ss_net_paid#1174,ss_net_paid_inc_tax#1175,ss_net_profit#1176] parquet ! : +- Distinct : : : +- Relation[d_date_sk#329,d_date_id#330,d_date#331,d_month_seq#332,d_week_seq#333,d_quarter_seq#334,d_year#335,d_dow#336,d_moy#337,d_dom#338,d_qoy#339,d_fy_year#340,d_fy_quarter_seq#341,d_fy_week_seq#342,d_day_name#343,d_quarter_name#344,d_holiday#345,d_weekend#346,d_following_holiday#347,d_first_dom#348,d_last_dom#349,d_same_day_ly#350,d_same_day_lq#351,d_current_day#352,... 4 more fields] parquet ! : +- Project [c_last_name#1952, c_first_name#1951, d_date#1917] : : +- Relation[c_customer_sk#154,c_customer_id#155,c_current_cdemo_sk#156,c_current_hdemo_sk#157,c_current_addr_sk#158,c_first_shipto_date_sk#159,c_first_sales_date_sk#160,c_salutation#161,c_first_name#162,c_last_name#163,c_preferred_cust_flag#164,c_birth_day#165,c_birth_month#166,c_birth_year#167,c_birth_country#168,c_login#169,c_email_address#170,c_last_review_date#171] parquet ! : +- Filter (((cs_sold_date_sk#872 = d_date_sk#1915) AND (cs_bill_customer_sk#875 = c_customer_sk#1943)) AND ((d_month_seq#1918 >= 1200) AND (d_month_seq#1918 <= (1200 + 11)))) : +- Distinct ! : +- Join Inner : +- Project [c_last_name#1952, c_first_name#1951, d_date#1917] ! : :- Join Inner : +- Filter (((cs_sold_date_sk#872 = d_date_sk#1915) AND (cs_bill_customer_sk#875 = c_customer_sk#1943)) AND ((d_month_seq#1918 >= 1200) AND (d_month_seq#1918 <= (1200 + 11)))) ! : : :- Relation[cs_sold_date_sk#872,cs_sold_time_sk#873,cs_ship_date_sk#874,cs_bill_customer_sk#875,cs_bill_cdemo_sk#876,cs_bill_hdemo_sk#877,cs_bill_addr_sk#878,cs_ship_customer_sk#879,cs_ship_cdemo_sk#880,cs_ship_hdemo_sk#881,cs_ship_addr_sk#882,cs_call_center_sk#883,cs_catalog_page_sk#884,cs_ship_mode_sk#885,cs_warehouse_sk#886,cs_item_sk#887,cs_promo_sk#888,cs_order_number#889,cs_quantity#890,cs_wholesale_cost#891,cs_list_price#892,cs_sales_price#893,cs_ext_discount_amt#894,cs_ext_sales_price#895,... 10 more fields] parquet : +- Join Inner ! : : +- Relation[d_date_sk#1915,d_date_id#1916,d_date#1917,d_month_seq#1918,d_week_seq#1919,d_quarter_seq#1920,d_year#1921,d_dow#1922,d_moy#1923,d_dom#1924,d_qoy#1925,d_fy_year#1926,d_fy_quarter_seq#1927,d_fy_week_seq#1928,d_day_name#1929,d_quarter_name#1930,d_holiday#1931,d_weekend#1932,d_following_holiday#1933,d_first_dom#1934,d_last_dom#1935,d_same_day_ly#1936,d_same_day_lq#1937,d_current_day#1938,... 4 more fields] parquet : :- Join Inner ! : +- Relation[c_customer_sk#1943,c_customer_id#1944,c_current_cdemo_sk#1945,c_current_hdemo_sk#1946,c_current_addr_sk#1947,c_first_shipto_date_sk#1948,c_first_sales_date_sk#1949,c_salutation#1950,c_first_name#1951,c_last_name#1952,c_preferred_cust_flag#1953,c_birth_day#1954,c_birth_month#1955,c_birth_year#1956,c_birth_country#1957,c_login#1958,c_email_address#1959,c_last_review_date#1960] parquet : : :- Relation[cs_sold_date_sk#872,cs_sold_time_sk#873,cs_ship_date_sk#874,cs_bill_customer_sk#875,cs_bill_cdemo_sk#876,cs_bill_hdemo_sk#877,cs_bill_addr_sk#878,cs_ship_customer_sk#879,cs_ship_cdemo_sk#880,cs_ship_hdemo_sk#881,cs_ship_addr_sk#882,cs_call_center_sk#883,cs_catalog_page_sk#884,cs_ship_mode_sk#885,cs_warehouse_sk#886,cs_item_sk#887,cs_promo_sk#888,cs_order_number#889,cs_quantity#890,cs_wholesale_cost#891,cs_list_price#892,cs_sales_price#893,cs_ext_discount_amt#894,cs_ext_sales_price#895,... 10 more fields] parquet ! +- Distinct : : +- Relation[d_date_sk#1915,d_date_id#1916,d_date#1917,d_month_seq#1918,d_week_seq#1919,d_quarter_seq#1920,d_year#1921,d_dow#1922,d_moy#1923,d_dom#1924,d_qoy#1925,d_fy_year#1926,d_fy_quarter_seq#1927,d_fy_week_seq#1928,d_day_name#1929,d_quarter_name#1930,d_holiday#1931,d_weekend#1932,d_following_holiday#1933,d_first_dom#1934,d_last_dom#1935,d_same_day_ly#1936,d_same_day_lq#1937,d_current_day#1938,... 4 more fields] parquet ! +- Project [c_last_name#1998, c_first_name#1997, d_date#1963] : +- Relation[c_customer_sk#1943,c_customer_id#1944,c_current_cdemo_sk#1945,c_current_hdemo_sk#1946,c_current_addr_sk#1947,c_first_shipto_date_sk#1948,c_first_sales_date_sk#1949,c_salutation#1950,c_first_name#1951,c_last_name#1952,c_preferred_cust_flag#1953,c_birth_day#1954,c_birth_month#1955,c_birth_year#1956,c_birth_country#1957,c_login#1958,c_email_address#1959,c_last_review_date#1960] parquet ! +- Filter (((ws_sold_date_sk#1013 = d_date_sk#1961) AND (ws_bill_customer_sk#1017 = c_customer_sk#1989)) AND ((d_month_seq#1964 >= 1200) AND (d_month_seq#1964 <= (1200 + 11)))) +- Distinct ! +- Join Inner +- Project [c_last_name#1998, c_first_name#1997, d_date#1963] ! :- Join Inner +- Filter (((ws_sold_date_sk#1013 = d_date_sk#1961) AND (ws_bill_customer_sk#1017 = c_customer_sk#1989)) AND ((d_month_seq#1964 >= 1200) AND (d_month_seq#1964 <= (1200 + 11)))) ! : :- Relation[ws_sold_date_sk#1013,ws_sold_time_sk#1014,ws_ship_date_sk#1015,ws_item_sk#1016,ws_bill_customer_sk#1017,ws_bill_cdemo_sk#1018,ws_bill_hdemo_sk#1019,ws_bill_addr_sk#1020,ws_ship_customer_sk#1021,ws_ship_cdemo_sk#1022,ws_ship_hdemo_sk#1023,ws_ship_addr_sk#1024,ws_web_page_sk#1025,ws_web_site_sk#1026,ws_ship_mode_sk#1027,ws_warehouse_sk#1028,ws_promo_sk#1029,ws_order_number#1030,ws_quantity#1031,ws_wholesale_cost#1032,ws_list_price#1033,ws_sales_price#1034,ws_ext_discount_amt#1035,ws_ext_sales_price#1036,... 10 more fields] parquet +- Join Inner ! : +- Relation[d_date_sk#1961,d_date_id#1962,d_date#1963,d_month_seq#1964,d_week_seq#1965,d_quarter_seq#1966,d_year#1967,d_dow#1968,d_moy#1969,d_dom#1970,d_qoy#1971,d_fy_year#1972,d_fy_quarter_seq#1973,d_fy_week_seq#1974,d_day_name#1975,d_quarter_name#1976,d_holiday#1977,d_weekend#1978,d_following_holiday#1979,d_first_dom#1980,d_last_dom#1981,d_same_day_ly#1982,d_same_day_lq#1983,d_current_day#1984,... 4 more fields] parquet :- Join Inner ! +- Relation[c_customer_sk#1989,c_customer_id#1990,c_current_cdemo_sk#1991,c_current_hdemo_sk#1992,c_current_addr_sk#1993,c_first_shipto_date_sk#1994,c_first_sales_date_sk#1995,c_salutation#1996,c_first_name#1997,c_last_name#1998,c_preferred_cust_flag#1999,c_birth_day#2000,c_birth_month#2001,c_birth_year#2002,c_birth_country#2003,c_login#2004,c_email_address#2005,c_last_review_date#2006] parquet : :- Relation[ws_sold_date_sk#1013,ws_sold_time_sk#1014,ws_ship_date_sk#1015,ws_item_sk#1016,ws_bill_customer_sk#1017,ws_bill_cdemo_sk#1018,ws_bill_hdemo_sk#1019,ws_bill_addr_sk#1020,ws_ship_customer_sk#1021,ws_ship_cdemo_sk#1022,ws_ship_hdemo_sk#1023,ws_ship_addr_sk#1024,ws_web_page_sk#1025,ws_web_site_sk#1026,ws_ship_mode_sk#1027,ws_warehouse_sk#1028,ws_promo_sk#1029,ws_order_number#1030,ws_quantity#1031,ws_wholesale_cost#1032,ws_list_price#1033,ws_sales_price#1034,ws_ext_discount_amt#1035,ws_ext_sales_price#1036,... 10 more fields] parquet ! : +- Relation[d_date_sk#1961,d_date_id#1962,d_date#1963,d_month_seq#1964,d_week_seq#1965,d_quarter_seq#1966,d_year#1967,d_dow#1968,d_moy#1969,d_dom#1970,d_qoy#1971,d_fy_year#1972,d_fy_quarter_seq#1973,d_fy_week_seq#1974,d_day_name#1975,d_quarter_name#1976,d_holiday#1977,d_weekend#1978,d_following_holiday#1979,d_first_dom#1980,d_last_dom#1981,d_same_day_ly#1982,d_same_day_lq#1983,d_current_day#1984,... 4 more fields] parquet ! ``` ReplaceIntersectWithSemiJoin 属于 **Replace Operators** 规则组,主要是用于将 `INTERSECT` 求两表交集的操作,替换为 `LEFT SEMI JOIN` 的形式,其源码如下: ```scala object ReplaceIntersectWithSemiJoin extends Rule[LogicalPlan] { def apply(plan: LogicalPlan): LogicalPlan = plan transform { case Intersect(left, right, false) => assert(left.output.size == right.output.size) val joinCond = left.output.zip(right.output).map { case (l, r) => EqualNullSafe(l, r) } Distinct(Join(left, right, LeftSemi, joinCond.reduceLeftOption(And), JoinHint.NONE)) } } ``` 替换流程如下: 1. 检查左右表的字段数量是否相同 2. 将左表与右表的字段属性进行 zip,并包装为 `EqualNullSafe`,后续会作为 Join 的关联条件 - EqualNullSafe 是一个二元的断言表达式,与 EqualTo 相比,在 NULL 方面的对比规则有所不同 3. 之后根据上一步得到的 Join 条件,构造一个 Join 节点,Join 类型为 LeftSemi 4. 在分布式环境下,虽然单个 Partition 满足了 INTERSECT,但不同 Partition 可能存在重复结果,因此 还需要添加一个 Distinct 节点,将各个节点满足 Join 条件的结果 Shuffle 到相同节点进行去重 ### 作业二:架构设计题 ![image-20211114151110902](https://gitee.com/zhxuankun/Image/raw/master/ARTS_Tips/image-20211114151110902.png) Lambda 主要分为实时处理与批处理两套。 批处理: 1. 数据通过 Flume 或者 Filebeat + Logstash 等方式存入 HDFS 集群中 2. 使用 AirFlow、Oozie 等调度工具,启动定时任务进行数据清理、转换,再次写入 HDFS,并以 Parquet 格式存储,根据具体业务需求构造数仓的各个存储层次 3. 基于 Hive 建表,并提供 JDBC 接口 4. Presto 与 SparkSQL 基于 HIve Metastore 的 JDBC 接口,小数据量使用 Presto 提供 ad-hoc 查询功能,大数据量选择 SparkSQL,由 Gateway 提供查询引擎的选择、失败降级等功能 流处理: 1. 数据通过 Flume 或者 Filebeat + Logstash 等方式送到 Kafka 集群中 2. 使用 Flink 实时消费,并写入 HBase 3. 基于 HBase 做一些实时方面的业务应用 Lambda 架构的优点: 1. 低延迟,流式处理层实时消费数据并处理,可以在一个较小的窗口内就将数据提供给用户使用 2. 数据一致性,Lambda 中数据都是线性处理的,不会出现分布式数据库中可能出现的不同节点查询到的值不一致的问题 3. 伸缩性,Lambda 架构没有限定死需要使用的技术架构,开发人员可以按照自身需求,替换技术栈 4. 容错性,即使流式处理层出现了故障,也可以通过批处理层的计算结果进行补充,同时批处理层即使出错了,也只需要进行重算即可 Lambda 架构的缺点: 1. 双倍的存储成本及计算成本 2. 需要维护两套不同的技术栈 3. 批流代码不同意,批处理与流处理可能存在计算结果、语义等方面的不同 4. 流式计算无法应对太复杂的计算要求,仍需要放到批处理 ### 作业三:简答题 简述 Spark Shuffle 的工作原理,要求不少于 300 字。 Shuffle 主要分为两个阶段:Shuffle Write -> Shuffle Read,Shuffle Write 时,数据会先写入一个 buffer,再落到磁盘上,目录的位置可以通过 spark.local.dir 指定。每个 Shuffle Write 结束后,都会返回一个 MapStatus 给 Driver,接下来的 Shuffle Read 再根据 MapStatus 获取需要拉取数据的位置。 Spark Shuffle 最原始的 Shuffle 方式为 **HashShuffle**,Map 端每个 task 都会进行数据落盘,并生成 Reduce task 个数的文件,总的文件数为 Map task * Reduce task。 为了避免大量文件随机读取造成的 I/O 问题,HashShuffle 后面优化为了一个 Spark Core 执行多个 Map task,同个 Core 上的 Map task 复用同个 buffer,最终会生成 Map core * Reduce task 个文件。 再到后面出现了 Sort-Based Shuffle 的方式,引入了排序。在数据进行落盘时,每个 task 会先在 buffer 内将数据按照 Partition 进行排序落盘,并生成多个有序的中间结果文件,在这之后会将这多个中间结果文件合并为一个数据文件及索引文件,也就是说生成的最终文件数为 2 * Map task 个。由于最后的文件是有序的,后续的 Reduce task 可以简单的根据索引中记录的偏移量直接一次读取所需的数据。 由于 Sort-Based Shuffle 引入了排序的开销,而在某些场景下并不需要引入排序,或者可以通过更加快速的方式实现,因此还存在两个切换方案: 1. bypass,每个 task 不再对数据进行排序,而是根据 Hash 将数据进行落盘,最终再合并为一个数据文件 + 一个索引文件,更加快速高效,但存在维护哈希表以及需要多个 buffer 的开销,因此其切换条件之一就是 Partition 需要小于 200,可通过 spark.shuffle.sort.bypassMergeThreshold 进行调整,同时要求 Map 端无聚合操作 2. Tungsten Sort,可以享受到 Tungsten 的优化,其触发条件如下: 1. Partition 大于等于 200,且小于 16777216(Tungsten 用 24 bit 指针记录分区号) 2. Map 端无聚合操作 3. 序列化时,单条记录不能大于 128 MB(Tungsten 用 27bit 指针存储数据偏移量) 4. Serializer 需要支持 relocation,如 Kryo