基于 Spark 4.2,分支
branch-4.2阅读时长约 10分钟 · 进阶
背景
Parser 只负责把 SQL 字符串变成一棵 Unresolved LogicalPlan:
'Project ['id]+- 'Filter ('id > 10) +- 'UnresolvedRelation [t]
这棵树里有很多“不知道”:
t是表?视图?临时视图?物理表?id是哪张表的列?类型是什么?>两边类型是否能比较?count是内置函数还是 UDF?
Analyzer 的工作就是把这些“不知道”变成确定的对象。
一、Analyzer 在执行链路中的位置
QueryExecution.scala:192:
private val lazyAnalyzed = LazyTry {val plan = executePhase(QueryPlanningTracker.ANALYSIS) {analyzer.executeAndCheck(sqlScriptExecuted, tracker)}tracker.setAnalyzed(plan)plan}
真正入口是 Analyzer.executeAndCheck,在 Analyzer.scala:331:
def executeAndCheck(plan: LogicalPlan, tracker: QueryPlanningTracker): LogicalPlan = {if (plan.analyzed) {plan} else {def runAnalysis(): LogicalPlan =HybridAnalyzer.fromLegacyAnalyzer(legacyAnalyzer = this, tracker = tracker).apply(plan)// ...}}
注意:Spark 4.2 已经引入了 HybridAnalyzer。老文章里常见的“Analyzer 直接执行 RuleExecutor”这句话不完整了。
二、HybridAnalyzer:4.x 的分析入口
HybridAnalyzer.scala:54:
classHybridAnalyzer(...) {入口在 HybridAnalyzer.scala:66:
def apply(plan: LogicalPlan): LogicalPlan = {if (singlePassResolverExtensions.isDefined) {resolveInSinglePass(plan)} else if (dualRunLegacyAndSinglePassResolver) {resolveInDualRun(plan)} else {resolveInFixedPoint(plan)}}
默认路径仍然是 fixed-point analyzer。HybridAnalyzer.scala:272-275:
private def resolveInFixedPoint(plan: LogicalPlan): LogicalPlan = {val resolvedPlan = legacyAnalyzer.executeAndTrack(plan, tracker)legacyAnalyzer.checkAnalysis(resolvedPlan)resolvedPlan}
所以当前可以这样理解:
QueryExecution.lazyAnalyzed → Analyzer.executeAndCheck → HybridAnalyzer.apply → resolveInFixedPoint → legacyAnalyzer.executeAndTrack → legacyAnalyzer.checkAnalysis
三、Analyzer 本质上是 RuleExecutor
Analyzer.scala:304:
class Analyzer(...) extends RuleExecutor[LogicalPlan] with CheckAnalysis with ...它继承了 RuleExecutor[LogicalPlan],所以 Analyzer 的工作方式就是:
定义一批 batches
每个 batch 里面放一组 rules
每个 rule 改写 LogicalPlan
batch 按策略执行一次或多次,直到 fixed point
Analyzer.scala:506:
override def batches: Seq[Batch] =earlyBatches ++ Seq(Batch("Resolution", fixedPoint, ...),Batch("Remove TempResolvedColumn", Once, ...),Batch("Post-Hoc Resolution", Once, ...),Batch("Cleanup", fixedPoint, ...))
最核心的是 Resolution batch,入口在 Analyzer.scala:507。
四、RuleExecutor:规则如何跑起来
RuleExecutor.scala:125:
abstract classRuleExecutor[TreeType <: TreeNode[_]] extendsLogging{策略定义:
case object Once extends Strategy // RuleExecutor.scala:150case class FixedPoint(...) extends Strategy // RuleExecutor.scala:156
Batch 定义:
case class Batch(name: String, strategy: Strategy, rules: Rule[TreeType]*) // RuleExecutor.scala:162执行主循环在 RuleExecutor.scala:215:
def execute(plan: TreeType): TreeType = {var curPlan = planbatches.foreach { batch =>var iteration = 1var continue = truewhile (continue) {val lastPlan = curPlancurPlan = batch.rules.foldLeft(curPlan) { case (plan, rule) =>rule(plan)}if (iteration >= batch.strategy.maxIterations || curPlan.fastEquals(lastPlan)) {continue = false}iteration += 1}}curPlan}
这段就是 fixed-point 的核心:
一轮规则跑完,如果计划没有变化,就停;否则继续跑下一轮。
五、Resolution batch 里有哪些规则
Analyzer.scala:507 开始的 Resolution batch 很长。几个最重要的规则:
| 规则 | 加入 batch 的位置 | 定义入口 | 作用 |
|---|---|---|---|
ResolveRelations | Analyzer.scala:510 | Analyzer.scala:1042 | 把表名解析成具体 relation |
ResolveReferences | Analyzer.scala:519 | Analyzer.scala:1495 | 把列名解析成 AttributeReference |
ResolveFunctions | Analyzer.scala:533 | Analyzer.scala:2282 | 把函数名解析成具体 Expression |
ResolveAliases | Analyzer.scala:540 | Analyzer.scala:670 | 给未命名表达式补 alias |
ResolveSubquery | Analyzer.scala:541 | Analyzer.scala:2501 | 解析子查询 |
这些规则的目标都是消灭 Unresolved*。
六、ResolveRelations:表名变 relation
Parser 产物里表名是:
'UnresolvedRelation [t]ResolveRelations 定义在 Analyzer.scala:1042。它会查:
临时视图
全局临时视图
session catalog
v2 catalog
最后把 UnresolvedRelation 替换成具体的逻辑节点,比如:
SubqueryAliasLogicalRelationDataSourceV2RelationView
这就是为什么“表不存在”是在 Analyzer 阶段报错,而不是 Parser 阶段。
七、ResolveReferences:列名变 AttributeReference
ResolveReferences 定义在 Analyzer.scala:1495,入口在 Analyzer.scala:1532。
它的工作是把:
'id变成:
id#0L这个 #0L 不是随便加的。它代表一个带 ExprId 和类型信息的 AttributeReference。
为什么需要 ExprId?因为 SQL 里可能有同名列:
SELECT id FROM t1 JOIN t2 ON t1.id = t2.id光看名字 id 不够,Analyzer 必须知道你指的是哪一个输出属性。
八、ResolveFunctions:函数名变 Expression
ResolveFunctions 定义在 Analyzer.scala:2282,入口在 Analyzer.scala:2283。
它把:
SELECT count(*), upper(name) FROM t里的 count、upper 解析成具体的 Catalyst expression:
聚合函数
普通标量函数
临时函数
catalog function
如果函数不存在,也是在 Analyzer 阶段报错。
九、checkAnalysis:不是解析完就完事
Analyzer 不只是“解析名字”。解析完还要检查计划是否合法。
CheckAnalysis.scala:306:
def checkAnalysis(plan: LogicalPlan): Unit = {val inlineCTE = InlineCTE(alwaysInline = true, keepDanglingRelations = true)val inlinedPlan = inlineCTE(plan)checkAnalysis0(inlinedPlan)plan.setAnalyzed()}
它会检查很多错误:
还有未解析的表 / 列 / 函数
聚合表达式是否合法
数据类型是否匹配
子查询是否合法
生成列 / 窗口函数 / streaming 限制等
CheckAnalysis.scala:327 的 plan.setAnalyzed() 很关键:只有通过检查的 plan 才会被标记为 analyzed。
此外,QueryExecution.assertSupported 在 QueryExecution.scala:161,会调用 UnsupportedOperationChecker.checkForBatch(analyzed)(QueryExecution.scala:163),检查一些批查询不支持的操作。
十、实战:观察 Analyzer 前后差异
spark.range(10).toDF("id").createOrReplaceTempView("t")val sql = "SELECT id, id + 1 AS next_id FROM t WHERE id > 3"val parsed = spark.sessionState.sqlParser.parsePlan(sql)val df = spark.sql(sql)println("=== Parsed ===")println(parsed)println("=== Analyzed ===")println(df.queryExecution.analyzed)
输出:
=== Parsed ==='Project ['id, ('id + 1) AS next_id#...]+- 'Filter ('id > 3)+- 'UnresolvedRelation [t], [], false=== Analyzed ===Project [id#0L, (id#0L + cast(1 as bigint)) AS next_id#...]+- Filter (id#0L > cast(3 as bigint))+- SubqueryAlias t+- View ... / Range ...
重点看三件事:
UnresolvedRelation消失'id变成id#0L1和3被补上类型转换
十一、总结
Analyzer 可以压缩成这条链:
Unresolved LogicalPlan→ HybridAnalyzer.apply→ legacyAnalyzer.executeAndTrack→ RuleExecutor batches→ Resolution batch→ ResolveRelations→ ResolveReferences→ ResolveFunctions→ ResolveAliases→ ResolveSubquery→ checkAnalysis→ Analyzed LogicalPlan
最值得带走的三个认知:
Analyzer 的核心目标是消灭 Unresolved。 表名、列名、函数名、子查询都在这里变成具体对象。
Analyzer 是 RuleExecutor。 规则按 batch 执行,FixedPoint batch 会反复迭代直到计划不再变化。
Spark 4.2 的入口是 HybridAnalyzer。 默认仍走 legacy fixed-point analyzer。
每天花费10分钟学习spark,让你技术之路走得更稳、更快。
喜欢的点个关注。
夜雨聆风