Spark大数据开发与应用案例(视频教学版)(十八)--第十二章
本书及作者信息
作者: 余辉 微信公众号:辉哥大数据
购买地址: 京东、 淘宝、 当当网
读者须知:本书配套示例源码、PPT课件、教学视频与作者答疑服务,购买之后可加粉丝群
本书封面

第12章 Spark SQL源码解读
本章将深入解读Spark SQL的源码精髓,带你领略从SQL执行到结果生成的全貌。我们将详细剖析Spark SQL的执行流程,从元数据管理的SessionCatalog出发,逐步探索SQL解析为逻辑执行计划、Analyzer绑定逻辑计划、Optimizer优化逻辑计划、SparkPlanner生成物理计划的关键步骤,并最终揭示如何从物理执行计划中获取inputRdd执行,全面解析Spark SQL的高效运行机制。
本章主要知识点:
- Spark SQL的执行过程
- 元数据管理SessionCatalog
- SQL解析成逻辑执行计划
- Analyzer绑定逻辑计划
- Optimizer优化逻辑计划
- 使用SparkPlanner生成物理计划
- 从物理执行计划获取inputRdd执行
12.1 Spark SQL的执行过程
正常的SQL执行先会经过SQL Parser解析SQL,然后经过Catalyst优化器处理,最后到Spark执行。而Catalyst的过程又分为几个过程,其中包括:
- SQL Parser:解析SQL语法、生成抽象语法树(AST)、解析校验语法。
- Analysis:利用Catalog信息将Unresolved Logical Plan解析成Analyzed Logical Plan来绑定元数据。
- Logical Optimizations:利用一些规则(Rule)将Analyzed Logical Plan解析成Optimized Logical Plan进行优化。
- Physical Planning:前面的Logical Plan不能被Spark执行,而这个过程是把Logical Plan转换成多个Physical Plans,然后利用代价模型(Cost Model)选择最佳的Physical Plan进行代码转换。
- Code Generation:这个过程会把SQL逻辑生成Java字节码,即生成代码。
- 分布式并行调度执行。
因此,整个SQL的执行过程可以使用图12-1表示。

图12-1中包括6个模块:Unresolved Logical Plan、Logical Plan、Optimized Logical Plan、Physical Plans、Cost Model和Selected Physical Plan,这6个模块就是Catalyst优化器处理的部分,也是本章讲解的主要内容。
12.2 元数据管理器SessionCatalog
SessionCatalog是Spark SQL的核心元数据管理器,负责会话级元数据的统一管控,涵盖数据库、表、视图、分区及函数等对象。它将临时视图/表注册至内存缓存,解析时通过Catalog路由元数据请求(如Hive表元数据由HiveCatalog从Hive Metastore获取)。在SQL解析阶段,Analyzer依赖SessionCatalog完成对象名称解析和元数据绑定,实现会话间元数据隔离与高效访问,是Spark SQL交互式查询性能优化的关键组件。
12.3 SQL解析成逻辑执行计划
当调用SparkSession的sql方法或者SQLContext的sql方法时,就会使用SparkSqlParser进行SQL解析。
Spark 2.0.0开始引入第三方语法解析器工具ANTLR,对SQL进行词法分析并构建语法树。Antlr是一款强大的语法生成器工具,可用于读取、处理、执行和翻译结构化的文本或二进制文件,是当前Java语言中使用最为广泛的语法生成器工具,常见的大数据SQL解析都用到了这个工具,包括Hive、Cassandra、Phoenix、Pig以及Presto等。目前新版本的Spark使用的是ANTLR4。
ANTLR4分为两个步骤来生成Unresolved LogicalPlan。
- 词法分析(SqlBaseLexer):Lexical Analysis负责将Token分组成符号类。
- 语法分析(SqlBaseParser):构建一棵分析树(Parse Tree)或者抽象语法树(Abstract Syntax Tree,AST)。
代码12-1是Scala语言编写的,定义了一个名为AstBuilder的类,这个类用于将ANTLR4生成的解析树(ParseTree)转换成催化剂(Catalyst)表达式、逻辑计划(LogicalPlan)或表标识符(TableIdentifier)。催化剂是Apache Spark SQL中用于查询优化和执行的核心组件。
代码12-1 AstBuilder.scala
// The AstBuilder converts an ANTLR4 ParseTree into a catalyst Expression,
// LogicalPlan or TableIdentifier.
class AstBuilder(conf: SQLConf) extends SqlBaseBaseVisitor[AnyRef] with Logging
{
import ParserUtils._
def this() = this(new SQLConf())
protected def typedVisit[T](ctx: ParseTree): T = {
...
}
}
具体来说,Spark基于Presto语法文件定义了Spark SQL语法文件SqlBase.g4(此文件位于 spark-2.4.3\sql\catalyst\src\main\antlr4\org\apache\spark\sql\catalyst\parser\SqlBase.g4),这个文件定义了Spark SQL支持的SQL语法,读者可以打开源代码阅读一下。
如果我们需要自定义新的语法,则需要在这个文件中定义好相关语法,然后使用ANTLR4对SqlBase.g4文件自动解析生成几个Java类,其中就包含重要的词法分析器SqlBaseLexer.java和语法分析器SqlBaseParser.java。运行SQL会使用SqlBaseLexer来解析关键词以及各种标识符等,然后使用SqlBaseParser来构建语法树。
下面以一条简单的SQL语句(见代码12-2)为例进行分析。
代码12-2 简单的SQL语句
SELECT sum(v)
FROM (
SELECT
t1.id,
1 + 2 + t1.value AS v
FROM t1 JOIN t2
WHERE
t1.id = t2.id AND
t1.cid = 1 AND
t1.did = t1.cid + 1 AND
t2.id > 5) o
SQL整个过程如图12-2所示。
生成语法树之后,使用AstBuilder将语法树转换成LogicalPlan,这个LogicalPlan也被称为 Unresolved LogicalPlan。解析后的逻辑计划如下:
== Parsed Logical Plan ==
'Project [unresolvedalias('sum('v), None)]
+- 'SubqueryAlias `bigdata_stu`
+- 'Project ['t1.id, ((1 + 2) + 't1.value) AS v#16]
+- 'Filter ((('t1.id = 't2.id) && ('t1.cid = 1)) && (('t1.did = ('t1.cid + 1)) && ('t2.id > 5)))
+- 'Join Inner
:- 'UnresolvedRelation `t1`
+- 'UnresolvedRelation `t2`
逻辑计划如图12-3所示。Unresolved LogicalPlan是从下往上看的,t1和t2两张表被生成了UnresolvedRelation,过滤的条件、选择的列以及聚合字段都知道了。Unresolved LogicalPlan仅仅是一种数据结构,不包含任何数据信息,比如不知道数据源、数据类型,以及不同的列来自哪张表等。
12.4 Analyzer绑定逻辑计划
Analyzer阶段会使用事先定义好的Rule以及SessionCatalog等信息对UnresolvedLogicalPlan进行元数据绑定。
代码12-3是用Scala语言编写的,用于Apache Spark SQL的查询分析阶段。它定义了两个主要的类Analyzer和SparkSqlParser。其中,Analyzer类负责逻辑查询计划的分析,SparkSqlParser类用于将SQL文本解析为逻辑计划。
代码12-3 Analyzer.scala
class Analyzer(
catalog: SessionCatalog,
conf: SQLConf,
maxIterations: Int)
extends RuleExecutor[LogicalPlan] with CheckAnalysis {
class SparkSqlParser(conf: SQLConf) extends AbstractSqlParser(conf) {
val astBuilder = new SparkSqlAstBuilder(conf)
override def parsePlan(sqlText: String): LogicalPlan = parse(sqlText) { parser =>
astBuilder.visitSingleStatement(parser.singleStatement()) match {
case plan: LogicalPlan => plan
case _ =>
val position = Origin(None, None)
throw new ParseException(Option(sqlText), "Unsupported SQL statement", position, position)
}
}
Rule定义在Analyzer中,具体如下:
lazy val batches: Seq[Batch] = Seq(
Batch("Hints", fixedPoint,
new ResolveHints.ResolveBroadcastHints(conf),
ResolveHints.ResolveCoalesceHints,
ResolveHints.RemoveAllHints),
Batch("Simple Sanity Check", Once,
LookupFunctions),
Batch("Substitution", fixedPoint,
CTESubstitution,
WindowsSubstitution,
EliminateUnions,
new SubstituteUnresolvedOrdinals(conf)),
Batch("Resolution", fixedPoint,
ResolveTableValuedFunctions :: // 解析表的函数
ResolveRelations :: // 解析表或视图
ResolveReferences :: // 解析列
ResolveCreateNamedStruct ::
ResolveDeserializer :: // 解析反序列化操作类
ResolveNewInstance ::
ResolveUpCast :: // 解析类型转换
ResolveGroupingAnalytics ::
ResolvePivot ::
ResolveOrdinalInOrderByAndGroupBy ::
ResolveAggAliasInGroupBy ::
ResolveMissingReferences ::
ExtractGenerator ::
ResolveGenerate ::
ResolveFunctions :: // 解析函数
ResolveAliases :: // 解析表别名
ResolveSubquery :: // 解析子查询
ResolveSubqueryColumnAliases ::
ResolveWindowOrder ::
ResolveWindowFrame ::
ResolveNaturalAndUsingJoin ::
ResolveOutputRelation ::
ExtractWindowExpressions ::
GlobalAggregates ::
ResolveAggregateFunctions ::
TimeWindowing ::
ResolveInlineTables(conf) ::
ResolveHigherOrderFunctions(catalog) ::
ResolveLambdaVariables(conf) ::
ResolveTimeZone(conf) ::
ResolveRandomSeed ::
TypeCoercion.typeCoercionRules(conf) ++
extendedResolutionRules : _*),
Batch("Post-Hoc Resolution", Once, postHocResolutionRules: _*),
Batch("View", Once,
AliasViewChild(conf)),
Batch("Nondeterministic", Once,
PullOutNondeterministic),
Batch("UDF", Once,
HandleNullInputsForUDF),
Batch("FixNullability", Once,
FixNullability),
Batch("Subquery", Once,
UpdateOuterReferences),
Batch("Cleanup", fixedPoint,
CleanupAliases)
)
从以上代码可以看出,多个性质类似的Rule组成一个Batch,而多个Batch构成一个Batches。这些Batches会由RuleExecutor执行,先按一个个Batch顺序执行,然后对Batch中的每个Rule顺序执行。每个Batch会执行一次(Once)或多次(FixedPoint,由spark.sql.optimizer.maxIterations参数决定),执行过程如图12-4所示。
12.5 Optimizer优化逻辑计划
Spark SQL优化器(Optimizer)中定义默认优化规则的核心逻辑(见图12-6),利用这些Rules对逻辑计划和Exepression进行迭代处理,从而使得树的节点进行合并和优化,在核心逻辑中包括三部分内容,分别介绍如下。
1)规则批次定义
defaultBatches方法返回一个优化规则批次序列(Seq[Batch]),包含基础操作符优化规则集operatorOptimizationRuleSet,涵盖投影下推、连接重排序、谓词下推等经典优化策略。
2)动态规则调整
实际执行时会通过excludedRules和nonExcludableRules进行规则过滤。
- excludedRules:需排除的规则列表。
- nonExcludableRules:强制保留的规则列表。
最终执行批次为defaultBatches - (excludedRules - nonExcludableRules),实现规则的灵活定制。
3)优化执行流程
规则按定义顺序依次作用于逻辑计划,通过模式匹配进行计划转换,代码如图12-5所示。例如下面的顺序:
- PushProjectionThroughUnion:将投影操作下推到UNION子节点。
- ReorderJoin:基于统计信息重排多表连接顺序。
- LimitPushDown:将LIMIT操作尽可能下推到扫描节点。

在前文的绑定逻辑计划阶段,对Unresolved LogicalPlan进行相关Transform操作得到了Analyzed Logical Plan。这个Analyzed Logical Plan可以直接转换成Physical Plan,然后在Spark中执行。然而,如果直接这么做,得到的Physical Plan很可能不是最优的,因为在实际应用中,很多低效的写法会带来执行效率的问题,需要进一步对Analyzed Logical Plan进行处理,得到更优的逻辑算子树。于是,针对SQL逻辑算子树的优化器Optimizer应运而生。
这个阶段的优化器主要是基于规则的优化器(Rule-based Optimizer,RBO),而绝大部分的规则都是启发式规则,也就是基于直观或经验而得出的规则,比如列裁剪(过滤掉查询不需要使用到的列)、谓词下推(将过滤尽可能下沉到数据源端)、常量累加(比如将1+2这种表达式事先计算好)以及常量替换(比如SELECT * FROM table WHERE i = 5 AND j = i + 3可以转换成SELECT * FROM table WHERE i = 5 AND j = 8)等。
与绑定逻辑计划阶段类似,这个阶段所有的规则也是通过实现Rule抽象类来定义的。多个规则组成一个Batch,多个Batch组成一个Batches,同样也是在RuleExecutor中执行。
RuleExecutor的核心源码骨架如图12-6~图12-10所示。




那么,针对代码12-2所示的SQL语句,其执行过程都会执行哪些优化呢?下面我们将举例说明。
12.5.1 谓词下推
谓词下推在Spark SQL中是由PushDownPredicate实现的,这个过程主要将过滤条件尽可能下推到底层,最好是数据源。对于代码12-2所示的SQL语句,使用谓词下推优化得到的逻辑计划如图12-11所示。
从图12-11可以看出,谓词下推将Filter算子直接下推到Join之前了(注意,图12-11是从下往上看的)。也就是在扫描t1表时会先使用((((isnotnull(cid#2) && isnotnull(did#3)) && (cid#2 = 1)) && (did#3 = 2)) && (id#0 > 5)) && isnotnull(id#0)过滤条件过滤出满足条件的数据。同时,在扫描t2表时,会先使用isnotnull(id#8) && (id#8 > 5)过滤条件过滤出满足条件的数据。经过这样的操作,可以大大减少Join算子处理的数据量,从而加快计算速度。
12.5.2 列裁剪
列裁剪在Spark SQL中是由ColumnPruning实现的。因为我们查询的表可能有很多个字段,但是对于每次查询,我们很大可能不需要扫描出所有的字段,这个时候利用列裁剪可以把那些查询不需要的字段过滤掉,使得扫描的数据量减少。因此,针对我们前面介绍的SQL,使用列裁剪优化得到的逻辑计划如图12-12所示。
从图12-12可以看出,经过列裁剪后,t1表只需要查询id和value两个字段,t2表只需要查询id字段。这样不仅减少了数据的传输,而且如果底层的文件格式为列存(比如Parquet),可以大大提高数据的扫描速度。
12.5.3 常量替换
常量替换在Spark SQL中是由ConstantPropagation实现的。其核心功能是将变量替换成常量,比如SELECT * FROM table WHERE i=5 AND j=i+3可以转换成SELECT * FROM table WHERE i=5 AND j=8。虽然这种优化在单条语句中看起来似乎并不显著,但如果处理的数据量非常大,涉及大量行的扫描,这种优化可以显著减少计算时间的开销。经过这种优化后得到的逻辑计划如图12-13所示。
我们的查询中有t1.cid=1 AND t1.did=t1.cid+1查询语句。从中可以看出,t1.cid已经是确定的值,所以我们完全可以使用它计算出t1.did的值。
12.5.4 常量累加
常量累加在Spark SQL中是由ConstantFolding实现的。这与常量替换类似,也是在优化阶段把一些常量表达式事先计算好。虽然这种优化在单个查询中看起来改动不大,但在处理大规模数据时,它可以显著减少计算量,从而降低CPU等资源的使用。经过这种优化后得到的逻辑计划如图12-14所示。
经过上述4个步骤的优化之后,得到的逻辑计划如下:
== Optimized Logical Plan ==
Aggregate [sum(cast(v#16 as bigint)) AS sum(v)#22L]
+- Project [(3 + value#1) AS v#16]
+- Join Inner, (id#0 = id#8)
:- Project [id#0, value#1]
: +- Filter (((((isnotnull(cid#2) && isnotnull(did#3)) && (cid#2 = 1)) && (did#3 = 2)) && (id#0 > 5)) && isnotnull(id#0))
: +- Relation[id#0,value#1,cid#2,did#3] csv
+- Project [id#8]
+- Filter (isnotnull(id#8) && (id#8 > 5))
+- Relation[id#8,value#9,cid#10,did#11] csv
最终得到的优化之后的逻辑计划如图12-15所示。
至此,优化逻辑计划阶段就完成了。
12.6 使用SparkPlanner生成物理计划
SparkSpanner使用Planning Strategies对优化后的逻辑计划进行转换,生成可以执行的物理计划SparkPlan。
/**
* 将逻辑计划转换成物理计划的抽象类
* 各实现类通过各种GenericStrategy来生成各种可行的待选物理计划
* 如果一个策略无法对逻辑计划树的所有操作进行转换,则会调用[GenericStrategy#planLater
* planLater]]来获得一个“占位符”对象暂时填充;之后由[[collectPlaceholders collected]]
* 收集并使用其他策略进行转换
* TODO: 到目前为止,永远只生成一个物理计划
* 后续迭代中会对“多计划”予以实现
*/
abstract class QueryPlanner[PhysicalPlan <: TreeNode[PhysicalPlan]] {
/** A list of execution strategies that can be used by the planner */
def strategies: Seq[GenericStrategy[PhysicalPlan]]
def plan(plan: LogicalPlan): Iterator[PhysicalPlan] = {
// 显然,此处还有大量工作需要做
// 收集所有可选的物理计划
val candidates = strategies.iterator.flatMap(_(plan))
abstract class SparkStrategies extends QueryPlanner[SparkPlan] {
self: SparkPlanner =>
// Plans special cases of limit operators
object SpecialLimits extends Strategy {
class SparkPlanner(
val sparkContext: SparkContext,
val conf: SQLConf,
val experimentalMethods: ExperimentalMethods)
extends SparkStrategies {
逻辑计划翻译成物理计划时,使用的是策略(Strategy)。前面介绍的逻辑计划绑定和优化经过Transformations动作之后,树的类型并没有改变。
Logical Plan转换成物理计划后,树的类型发生了改变,由Logical Plan转换成Physical Plan了。一个逻辑计划(Logical Plan)经过一系列的策略处理之后,得到多个物理计划(Physical Plans),物理计划在Spark中是由SparkPlan实现的。多个物理计划经过代价模型(Cost Model)得到选择后的物理计划(Selected Physical Plan),整个过程如图12-16所示。
Cost Model对应的是基于代价的优化(Cost-based Optimizations,CBO)。CBO主要由华为团队实现(详见SPARK-16026)。其核心思想是计算每个物理计划的代价,然后得到最优的物理计划。目前,这一部分并没有实现,直接返回多个物理计划列表的第一个作为最优的物理计划,代码如下:
lazy val sparkPlan: SparkPlan = {
SparkSession.setActiveSession(sparkSession)
// TODO: We use next(), i.e. take the first plan returned by the planner, here for now, but we will implement to choose the best plan.
planner.plan(ReturnAnswer(optimizedPlan)).next()
}
而SPARK-16026引入的CBO优化,主要是在前面介绍的优化逻辑计划阶段(Optimizer阶段)进行的,对应的Rule为CostBasedJoinReorder,并且默认是关闭的,需要通过spark.sql.cbo.enabled 或 spark.sql.cbo.joinReorder.enabled参数开启。
因此,到了这个节点,得到的物理计划如下:
== Physical Plan ==
*(3) HashAggregate(keys=[], functions=[sum(cast(v#16 as bigint))], output=[sum(v)#22L])
+- Exchange SinglePartition
+- *(2) HashAggregate(keys=[], functions=[partial_sum(cast(v#16 as bigint))], output=[sum#24L])
+- *(2) Project [(3 + value#1) AS v#16]
+- *(2) BroadcastHashJoin [id#0], [id#8], Inner, BuildRight
:- *(2) Project [id#0, value#1]
: +- *(2) Filter (((((isnotnull(cid#2) && isnotnull(did#3)) && (cid#2 = 1)) && (did#3 = 2)) && (id#0 > 5)) && isnotnull(id#0))
: +- *(2) FileScan csv [id#0,value#1,cid#2,did#3] Batched: false, Format: CSV, Location: InMemoryFileIndex[file:/iteblog/t1.csv], PartitionFilters: [], PushedFilters: [IsNotNull(cid), IsNotNull(did), EqualTo(cid,1), EqualTo(did,2), GreaterThan(id,5), IsNotNull(id)], ReadSchema: struct<id:int,value:int,cid:int,did:int>
+- BroadcastExchange HashedRelationBroadcastMode(List(cast(input[0, int, true] as bigint)))
+- *(1) Project [id#8]
+- *(1) Filter (isnotnull(id#8) && (id#8 > 5))
+- *(1) FileScan csv [id#8] Batched: false, Format: CSV, Location: InMemoryFileIndex[file:/iteblog/t2.csv], PartitionFilters: [], PushedFilters: [IsNotNull(id), GreaterThan(id,5)], ReadSchema: struct<id:int>
从上面的结果可以看出,物理计划阶段已经知道数据源是从CSV文件中读取的了,也已经知道文件的路径、数据类型等信息。而且在读取文件时,直接将过滤条件(PushedFilters)加进去了。
同时,这个Join变成了BroadcastHashJoin,即把t2表的数据广播到t1表所在的节点,如图12-17所示。
至此,Physical Plan就完全生成了。
12.7 从物理执行计划获取inputRdd执行
从物理计划上可以获取inputRdd。从物理计划上生成全阶段代码,并编译反射出迭代器newBiIterator的类[类名:BufferedRowIterator]。然后,对inputRDD进行一个转换(Transformation),得到最终要执行的RDD。
inputRdd.mapPartitionsWithIndex((index,iter)=>{
new newBiIterator(){
hasNext(){
iter.hasNext
}
next(){
processNext(iter.next())
}
}
})
然后,对最后返回的RDD执行所需要的行动算子:
rdd.collect().foreach(println)
12.8 本章小结
本章全面深入地解读了Spark SQL的源码,为我们揭示了Spark SQL从接收到SQL查询到最终执行并返回结果的完整流程。首先,我们了解了Spark SQL的执行过程,这是理解后续步骤的基础。接着,我们深入探讨了元数据管理器SessionCatalog的作用,它负责管理和维护数据库中的元数据。在SQL解析阶段,我们学习了如何将SQL查询解析为逻辑执行计划。随后,Analyzer模块将逻辑计划与数据表的元数据绑定,为执行计划提供了必要的信息。Optimizer模块则对逻辑计划进行优化,以提高查询的执行效率。优化后的逻辑计划通过SparkPlanner生成物理计划,这是执行查询的具体步骤。最后,我们从物理执行计划中获取inputRdd并执行,从而得到查询结果。这一系列步骤共同构成了Spark SQL高效、强大的查询处理能力。通过本章的学习,我们对Spark SQL的工作机制有了更深入的了解和认识。
本书其他章节
- Spark大数据开发与应用案例(视频教学版)(一)–文前
- Spark大数据开发与应用案例(视频教学版)(二)–第一章上
- Spark大数据开发与应用案例(视频教学版)(三)–第一章下
- Spark大数据开发与应用案例(视频教学版)(四)–第二章上
- Spark大数据开发与应用案例(视频教学版)(五)–第二章下
- Spark大数据开发与应用案例(视频教学版)(六)–第三章上
- Spark大数据开发与应用案例(视频教学版)(七)–第三章下
- Spark大数据开发与应用案例(视频教学版)(八)–第四章上
- Spark大数据开发与应用案例(视频教学版)(九)–第四章下
- Spark大数据开发与应用案例(视频教学版)(十)–第五章
- Spark大数据开发与应用案例(视频教学版)(十一)–第六章
- Spark大数据开发与应用案例(视频教学版)(十二)–第七章
- Spark大数据开发与应用案例(视频教学版)(十三)–第八章
- Spark大数据开发与应用案例(视频教学版)(十四)–第九章
- Spark大数据开发与应用案例(视频教学版)(十五)–第十章上
- Spark大数据开发与应用案例(视频教学版)(十六)–第十章下
- Spark大数据开发与应用案例(视频教学版)(十七)–第十一章
- Spark大数据开发与应用案例(视频教学版)(十八)–第十二章
- Spark大数据开发与应用案例(视频教学版)(十九)–第十三章
- Spark大数据开发与应用案例(视频教学版)(二十)–第十四章
- Spark大数据开发与应用案例(视频教学版)(二十一)–第十五章

更多推荐



所有评论(0)