# jike-task-06 **Repository Path**: osky1993/jike-task-06 ## Basic Information - **Project Name**: jike-task-06 - **Description**: Spark SQL作业 - **Primary Language**: Java - **License**: MIT - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2021-09-04 - **Last Updated**: 2022-07-04 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # jike-task-06 #### 介绍 Spark SQL作业 1. 为Spark SQL添加一条自定义命令 2. 构建SQL满足要求 3. 实现自定义优化规则(静默规则) 具体作业要求见同目录的PDF文件《Spark SQL作业-0829 课后作业》 ## 作业01-为Spark SQL添加一条自定义命令 #### 操作步骤 1. 去Github下载Spark源码(以前就下载好了的,更新就下就行),分支使用"branch-3.2" 2. build spark(进入spark源码目录,执行"./build/sbt package -DskipTests -Phive -Phive-thriftserver") 3. 修改SqlBase.g4,具体修改内容见下 4. 使用antlr4插件解析SqlBase.g4 5. 增加spark version命令 6. 再次build spark 7. 执行"bin/spark-sql -S",输入"show version"进行验证 上述需要修改或者调整的文件见"work01"目录 #### SqlBase.g4文件的修改 ```antlrv4 1. 增加statement statement : query | SHOW VERSION #showVersion 2. 增加keyword //============================ // Start of the keywords list //============================ //--SPARK-KEYWORD-LIST-START VERSION: 'VERSION'; 3.标识为非保留 3.1 nonReserved //--DEFAULT-NON-RESERVED-START : ADD | VERSION 3.2 ansiNonReserved //--ANSI-NON-RESERVED-START : ADD | VERSION ``` #### 使用antlr4解析SqlBase.g4 双击执行即可,必须使用*V*P*N*,不然会让你怀疑人生的,还有maven优先选择中央仓库,aliyun的仓库暗坑不少!!! ![](https://i.loli.net/2021/09/05/6ORBNDT8YkImCgS.png) #### 增加spark version命令 修改sparkSqlParser.scala文件, 重写visitShowVersion方法 ```scala override def visitShowVersion(ctx: ShowVersionContext): LogicalPlan = withOrigin(ctx) { ShowVersionCommand() } ``` 增加ShowVersionCommand实现类 ```scala package org.apache.spark.sql.execution.command import org.apache.spark.sql.{Row, SparkSession} import org.apache.spark.sql.catalyst.expressions.{Attribute, AttributeReference} import org.apache.spark.sql.types.StringType case class ShowVersionCommand() extends LeafRunnableCommand { override def output: Seq[Attribute] = Seq(AttributeReference("version", StringType, nullable = true)()) override def run(sparkSession: SparkSession): Seq[Row] = { val javaVersion = System.getProperty("java.version") .split("[+.\\-]", 3) .mkString("Array(", ", ", ")") val sparkVersion = sparkSession.version Seq(Row("Spark version: " + sparkVersion + " Java version:" + javaVersion)) } } ``` #### 结果展示 ![](https://i.loli.net/2021/09/05/dTJCBMmvz82sOYH.png) ## 作业02-构建SQL满足要求 ### 题目一 #### 要求 构建一条SQL,同时apply下面三条优化规则: + CombineFilters + CollapseProject + BooleanSimplification 提前准备,设置日志级别和建表 ```sql set spark.sql.planChangeLog.level=WARN; create table students (name string, age int, student_num string); ``` #### SQL如下 ```sql select a.student_num from ( select name, student_num, age from students where 1=1 and age > 5 ) a where a.age<30; ``` 运行结果日志见"work02/spark-01.log" #### 分析 + CombineFilters: 相邻(上下级)Filter操作合并,把从表达式里面的过滤替换成,先做过滤再取表达式,并且掉过滤里面的别名属性 + CollapseProject: 该规则用于删除不必要的projects(投影) + BooleanSimplification: 提前短路可以短路的布尔表达式, 简化filter,比如where 1=1 或者where 1=2,前者直接去掉这个过滤,后者这个查询就没必要做了 ### 题目二 构建一条SQL,同时apply下面五条优化规则: + ConstantFolding + PushDownPredicates + ReplaceDistinctWithAggregate + ReplaceExceptWithAntiJoin + FoldablePropagation #### SQL如下 ```sql (select a.student_num , a.age + (5 + 10) , Now() nowtime from ( select distinct name, age , student_num from students ) a where a.age>20 order by nowtime) except ( select a.student_num , a.age + (5 + 10), Now() nowtime from ( select distinct name, age , student_num from students ) a where a.name="osky"); ``` 运行结果日志见"work02/spark-02.log" #### 分析 + ConstantFolding: 常量折叠 + PushDownPredicates: 谓词下推,把过滤算子(就是在sql语句里面写的where语句),尽可能地放在执行计划靠前的地方 + ReplaceDistinctWithAggregate: 将 distinct 操作转为 aggregate + ReplaceExceptWithAntiJoin: 将 except 操作转化为 anti join + FoldablePropagation: 可折叠的传递 ## 作业03-实现自定义优化规则(静默规则) //TODO