spark1: scala for spark
第 01 章 · Scala 速译:用 Java 读懂 Spark 源码
Spark 源码精读系列 · 第 01 篇 · 地基 本篇不教 Scala,只教"把 Scala 翻译成你脑子里的 Java"。 目标不是让你会写 Scala,而是让你看到 Spark 源码里任意一段代码时,能准确说出它在干什么。
读前须知
前置知识 / 能力(请逐条自查)
本篇的类比基准全部是 Java,以下是你需要确认自己具备的:
- Java 泛型:能读懂
List<String>、Map<K, V>;知道<T extends Number>是类型上界约束。 - Java 接口与抽象类:知道
interface只能定义方法签名(Java 8 后可有default方法)、abstract class可以有字段和构造函数、一个类只能extends一个父类但能implements多个接口。 - Java 的
Optional:知道它是用来避免null判空的容器类,有isPresent()/get()/orElse()。如果你没用过Optional,只需知道"它是个盒子,里面可能有值也可能没有"即可,本篇会重新讲。 - Java 的
equals/hashCode契约:知道两个对象"相等"要靠这两个方法定义,以及 IDE 或 Lombok 可以自动生成它们。 - Java 8 Lambda 与 Stream:能读懂
list.stream().filter(x -> x > 0).map(x -> x * 2)。这是本篇最重要的前置项——Scala 的集合操作和它几乎一一对应。 - 不可变对象概念:知道
String是不可变的,"修改"一个String实际上是产生新对象。
你不需要具备:任何 Scala 经验、函数式编程理论、类型系统理论。
读完你将获得
- 看到
case class Filter(condition: Expression, child: LogicalPlan) extends OrderPreservingUnaryNode with PredicateHelper这样一行,能准确说出它对应的 Java 结构是什么。 - 看到
def runJob[T, U: ClassTag](rdd: RDD[T], func: (TaskContext, Iterator[T]) => U, ...)这样一个长签名,能逐段拆出每个符号的含义——包括那个看起来最劝退的U: ClassTag。 - 能读懂 Spark 优化器规则的标准写法——认出
case ... =>是在匹配什么树形状,这是后续 17 章的核心阅读技能。 - 掌握 Scala 里
_的常见含义(它在 Spark 源码里出现频率极高,且含义随上下文变化)。 - 建立一张**「Scala 语法 → Java 等价物」速查表**,后续章节遇到不认识的写法可以回来查。
- 知道自己可以跳过哪些 Scala 特性(有些特性 Spark 用得极少,不值得花时间)。
- 拿到本章的 面试问答(带真实出现频率标记)——顺带知道为什么 Scala 语法本身在面试里几乎不问,以及本章的哪一条是真正值得准备的。
本篇的读法建议
不要试图记住。这一章的定位是词典,不是课文。第一遍通读建立印象,后面 17 章遇到读不懂的代码,回来查对应小节。
一、来龙去脉:为什么 Spark 是 Scala 写的,这对你意味着什么
为什么不是 Java
Spark 2009 年诞生于 UC Berkeley AMPLab。当时选 Scala 有几个现实原因:
- Scala 运行在 JVM 上,能直接复用整个 Hadoop / Java 生态(HDFS 客户端、序列化库等)。这是硬约束——大数据生态的既有资产全在 JVM 上。
- Scala 有 REPL(交互式解释器)。
spark-shell能存在,是因为 Scala 自带 REPL。2010 年前后的 Java 没有(Java 9 才有jshell),而"交互式探索数据"是 Spark 早期最重要的卖点之一。 - 函数式语法让 API 简洁。
rdd.map(x => x * 2)在 2010 年的 Java 6/7 里要写成一个匿名内部类,五行起步——Java 8 的 Lambda 是 2014 年才有的。
代价(这条对你很重要):Scala 的语法灵活度极高,同一件事有多种写法,且大量使用符号和隐式规则。这让读源码的门槛显著高于 Java。你现在感受到的困难是真实的、普遍的,不是你的问题。
一个让你安心的事实
Spark 源码只用了 Scala 的一个子集。 那些让 Scala 声名狼藉的高级特性——高阶类型(higher-kinded types)、类型类(type class)推导、宏、复杂的隐式解析链——在 Spark 的 SQL 模块里几乎不出现,或只出现在极少数底层工具类里。
你真正需要掌握的,是本篇列出的这十几个语法点。 掌握它们,Spark SQL 侧 90% 以上的代码你都能读懂。
二、第一课:读懂一行定义
我们直接从 Spark 的真实代码开始。这是逻辑计划里的"过滤"节点,也就是你 SQL 里 WHERE 子句最终变成的东西:
文件:sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/basicLogicalOperators.scala:335
case class Filter(condition: Expression, child: LogicalPlan)
extends OrderPreservingUnaryNode with PredicateHelper {
这一行半包含了 5 个 Scala 语法点。逐个拆。
2.1 类型在后面:name: Type 而不是 Type name
Scala 的变量和参数声明,类型写在冒号后面:
condition: Expression // Scala
Expression condition // Java 等价
为什么反过来:这是 Pascal/ML 系语言的传统,好处是类型可以省略时(类型推断)读起来仍然通顺。你只需要习惯它,没有更深的道理。
方法返回值同理,写在参数列表后面:
def output: Seq[Attribute] = child.output
List<Attribute> output() { return child.output(); } // 近似
2.2 [] 是泛型,不是数组
这是 Java 程序员读 Scala 最容易出错的地方:
Seq[Attribute] // Scala:一个元素类型为 Attribute 的序列
List<Attribute> // Java 等价
Scala 用方括号 [] 表示泛型(Java 用尖括号 <>)。而 Scala 里的数组写作 Array[Int]——注意它自己也是用方括号表示泛型的。
速记:在 Scala 里看到
[...],一律先当成 Java 的<...>来读。
2.3 case class:Java 的 record / Lombok @Data
case class Filter(condition: Expression, child: LogicalPlan)
case class 是 Scala 的一个关键字,它自动帮你生成一堆样板代码。最接近的 Java 类比是 Java 14+ 的 record,或者 Lombok 的 @Data + @AllArgsConstructor:
// Java 等价(用 record 表达)
public record Filter(Expression condition, LogicalPlan child) { }
// 或者用 Lombok 表达
@Data
@AllArgsConstructor
public class Filter {
private final Expression condition;
private final LogicalPlan child;
}
case class 自动提供的东西:
| 自动生成 | Java 里对应什么 | 在 Spark 里为什么重要 |
|---|---|---|
| 构造函数 | new Filter(cond, child) |
但 Scala 里可以省略 new,直接写 Filter(cond, child) |
所有参数变成 public final 字段 |
getCondition() 等 getter |
直接 filter.condition 访问,无需 getter |
equals / hashCode |
Lombok @EqualsAndHashCode |
按值比较——两个 Filter 只要 condition 和 child 相等就相等 |
toString |
Lombok @ToString |
explain 打印计划树时用的就是它 |
copy 方法 |
无直接对应 | 极其重要,见 2.6 |
| 可用于模式匹配 | 无直接对应 | 这是 Catalyst 的命脉,见第四节 |
为什么 Spark 的计划节点全是
case class:因为 Catalyst 优化器的工作方式是"匹配某种树形状 → 替换成另一种形状"。模式匹配和按值比较这两个能力是刚需,而case class白送这两个能力。这个设计动机会在第 04 章(TreeNode)完整展开。
2.4 extends ... with ...:单继承 + 多混入
extends OrderPreservingUnaryNode with PredicateHelper
extends OrderPreservingUnaryNode implements PredicateHelper // 近似 Java
规则:第一个用 extends,后面的都用 with。至于 extends 后面跟的是类还是 trait,语法上不区分。
2.5 trait:带实现的接口
PredicateHelper 是一个 trait。
定义:trait 是 Scala 的接口,但它可以包含具体方法实现和字段。
最接近的 Java 类比是 Java 8 的 interface + default 方法:
// Scala
trait PredicateHelper {
def splitConjunctivePredicates(condition: Expression): Seq[Expression] = {
// ... 具体实现
}
}
// Java 8+ 等价
interface PredicateHelper {
default List<Expression> splitConjunctivePredicates(Expression condition) {
// ... 具体实现
}
}
trait 比 Java interface 强的地方(这也是 Spark 大量用它的原因):
- 可以有字段(Java interface 只能有
public static final常量) - 可以有构造逻辑
- 一个类可以混入任意多个 trait
在 Spark 里 trait 的典型用途是"能力包":PredicateHelper 提供一组处理谓词(也就是 WHERE 条件)的工具方法,谁需要谁 with 一下。这和 Java 里"把工具方法放在 XxxUtils 静态类里"是解决同一个问题的两种方案——Scala 的做法是把工具方法直接混进类里,调用时不用写类名前缀。
排障提示:当你在 Spark 源码里看到一个方法调用,却在当前类里找不到它的定义时——第一反应应该是去看它
with了哪些 trait。这是读 Spark 源码时最常见的"方法找不到"的原因。
2.6 现在,完整读一遍
case class Filter(condition: Expression, child: LogicalPlan)
extends OrderPreservingUnaryNode with PredicateHelper {
中文翻译:定义一个名为 Filter 的不可变数据类,它有两个字段——condition(类型 Expression,表示过滤条件)和 child(类型 LogicalPlan,表示它的子节点);它继承自 OrderPreservingUnaryNode(一个"保持顺序的单子节点"基类),并混入 PredicateHelper 这个工具 trait。
Java 等价:
public record Filter(Expression condition, LogicalPlan child)
extends OrderPreservingUnaryNode implements PredicateHelper { }
// 注:Java record 实际不能 extends 类,这里只为表达语义
对应到你写的 SQL:WHERE o.dt = '2026-08-01' 这个条件,最终就变成一个 Filter 节点,其中 condition 字段装着 o.dt = '2026-08-01' 这个表达式,child 字段指向"读 orders 表"那个节点。
三、val / var / def / lazy val:四种成员
继续读 Filter 的方法体。文件同上,336–358 行:
override def output: Seq[Attribute] = child.output
final override val nodePatterns: Seq[TreePattern] = Seq(FILTER)
override protected lazy val validConstraints: ExpressionSet = {
...
}
先扫清三个陌生名词(本章只需知道它们"是什么",用途在后续章节展开):
Attribute:一个"列的引用"。SELECT city里的city就是一个Attribute。TreePattern:一个枚举值,标记"这个节点是什么类型"(这里是FILTER)。Catalyst 用它做快速剪枝——见 8.4 节。ExpressionSet:一个"表达式的集合",但它按语义去重(a > 1和1 < a会被认为是同一个)。validConstraints表示"经过这个节点后,哪些条件必然成立",供优化器推导用。
这里出现了三种不同的成员定义方式。加上 var,Scala 一共四种:
| Scala | 含义 | Java 等价 | 求值时机 |
|---|---|---|---|
val x = ... |
不可变值 | final 字段 |
对象构造时立即求值一次 |
var x = ... |
可变变量 | 普通字段 | 构造时立即求值 |
def x = ... |
方法 | 方法 | 每次调用都重新求值 |
lazy val x = ... |
延迟初始化的不可变值 | 需手写双重检查锁的单例字段 | 首次访问时求值一次,之后缓存 |
为什么这个区分在 Spark 里很关键
def 每次调用都重算。看这一行:
override def output: Seq[Attribute] = child.output
output 表示"这个节点输出哪些列"。它用 def 定义,意味着每次访问 filter.output 都会重新去问 child.output。对于一棵深度为 N 的计划树,访问根节点的 output 会递归到叶子,代价是 O(N)。
lazy val 只算一次:
override protected lazy val validConstraints: ExpressionSet = { ... }
validConstraints 的计算涉及集合运算(见下面代码),比较贵。用 lazy val 意味着:第一次访问时才计算,算完缓存起来,后续访问直接返回缓存值。
Java 里要实现同样的效果,你得手写:
// Java 等价:lazy val 的手工实现
private volatile ExpressionSet validConstraints;
public ExpressionSet getValidConstraints() {
if (validConstraints == null) { // 第一次检查
synchronized (this) {
if (validConstraints == null) { // 第二次检查(双重检查锁)
validConstraints = computeConstraints();
}
}
}
return validConstraints;
}
Scala 的 lazy val 就是这段代码的语法糖——编译器自动生成线程安全的延迟初始化。
为什么 Spark 大量使用
lazy val:Catalyst 的计划节点在优化过程中会被反复访问、反复重建。很多属性(输出列、约束条件、统计信息)计算代价不小,但在一个节点的生命周期内不会变。lazy val恰好匹配这个场景:贵、不变、可能根本不被访问。代价:
lazy val有隐藏成本——每次访问都要检查"是否已初始化",且早期 Scala 版本的实现涉及加锁。在极热的代码路径上这是可测量的开销。这也是为什么有些字段 Spark 明确用val而非lazy val。
override 与 final
final override val nodePatterns: Seq[TreePattern] = Seq(FILTER)
override:Scala 里覆写父类成员必须显式写override,不写编译报错。这比 Java 的@Override(可选注解)严格。这对你是好事——看到override就知道父类里一定有同名成员,可以去父类查它的文档。final:和 Java 一样,禁止子类再覆写。
注意 val 可以覆写 def:父类用 def output 定义,子类可以用 val output 覆写。这在 Java 里不存在(方法和字段是两个命名空间)。读源码时如果在父类看到 def、子类看到 val,它们是同一个成员。
四、模式匹配:Catalyst 的命脉
这是本章最重要的一节。 如果只能记住一件事,记这个。
4.1 从 switch 说起
Scala 的 match 表面上像 Java 的 switch:
// Scala
val result = x match {
case 1 => "one"
case 2 => "two"
case _ => "other"
}
// Java 14+ switch 表达式,几乎一一对应
var result = switch (x) {
case 1 -> "one";
case 2 -> "two";
default -> "other";
};
对应关系:
match≈switchcase X => ...≈case X -> ...case _ => ...≈default ->(_是通配符)
但 Scala 的 match 强大得多:它不只能匹配值,还能匹配对象的类型和结构,并同时把内部字段解构出来。
4.2 解构:匹配的同时拆开对象
看 Filter 里的这段(basicLogicalOperators.scala:339-342):
override def maxRows: Option[Long] = condition match {
case Literal.FalseLiteral => Some(0L)
case _ => child.maxRows
}
逐行翻译:
override def maxRows: Option[Long] = condition match {
// ^方法名 ^返回类型 ^对 condition 这个字段做匹配
case Literal.FalseLiteral => Some(0L)
// ^如果 condition 恰好是"常量 false" ^那么最大行数是 0
case _ => child.maxRows
// ^其他任何情况 ^最大行数取决于子节点
}
这段代码的业务含义(这才是重点):如果你写了 WHERE false,Spark 知道这个查询一行都不会返回,maxRows 直接是 0。这个信息后续会被优化器利用——既然最多返回 0 行,整棵子树可以直接删掉换成空结果。
Java 等价(用 instanceof 模拟):
Optional<Long> maxRows() {
if (condition.equals(Literal.FALSE_LITERAL)) {
return Optional.of(0L);
} else {
return child.maxRows();
}
}
4.3 真正的威力:解构嵌套结构
上面还只是"匹配值"。Catalyst 真正依赖的是匹配树的形状并拆出内部字段。
看这个真实的优化规则(Optimizer.scala:2037):
case Filter(fc, nf @ Filter(nc, grandChild)) if nc.deterministic =>
这一行信息量极大,逐个符号拆解:
case Filter(fc, nf @ Filter(nc, grandChild)) if nc.deterministic =>
// ^^^^^^ 匹配一个 Filter 节点
// ^^ 把它的 condition 字段命名为 fc
// ^^ 把整个内层节点命名为 nf
// ^ @ = "并且把它叫做"
// ^^^^^^ 内层必须也是一个 Filter
// ^^ 内层的 condition 命名为 nc
// ^^^^^^^^^^ 内层的 child 命名为 grandChild
// ^^^^^^^^^^^^^^^^^^^^ 守卫条件:仅当 nc 是确定性的
中文翻译:如果当前节点是一个 Filter,并且它的子节点也是一个 Filter,并且内层的过滤条件是确定性的(每次求值结果相同),那么就匹配成功;同时把外层条件叫 fc、内层整体叫 nf、内层条件叫 nc、内层的子节点叫 grandChild。
对应的 SQL 场景:
SELECT * FROM (SELECT * FROM orders WHERE dt = '2026-08-01') t WHERE amount > 100
这会产生两个嵌套的 Filter 节点。这条规则(CombineFilters)就是要把它们合并成一个 Filter,条件用 AND 连起来,避免数据被遍历两遍。
Java 等价(你会立刻看出 Scala 的价值):
// Java 等价:手写要这么长
if (plan instanceof Filter) {
Filter outer = (Filter) plan;
Expression fc = outer.condition();
LogicalPlan innerPlan = outer.child();
if (innerPlan instanceof Filter) {
Filter nf = (Filter) innerPlan;
Expression nc = nf.condition();
LogicalPlan grandChild = nf.child();
if (nc.deterministic()) {
// ... 终于可以干活了
}
}
}
Scala 一行 = Java 十行。 而 Catalyst 有 200+ 条这样的规则。这就是为什么 Spark 选 Scala 写优化器——模式匹配让"匹配树形状并改写"这件事的代码量降了一个数量级。
4.4 守卫(guard):if 跟在 case 后面
case Filter(fc, nf @ Filter(nc, grandChild)) if nc.deterministic =>
// ^^^^^^^^^^^^^^^^^^ 守卫
case 模式 if 条件 => 表示:形状匹配上了,还要额外满足这个条件才算数。
这个具体的守卫为什么存在(业务含义,非常典型的排障知识):deterministic 表示"确定性"——同样的输入永远得到同样的输出。像 rand()、current_timestamp() 这类函数是非确定性的。
如果内层条件是 rand() < 0.1(随机采样 10%),把两个 Filter 合并会改变求值次数和语义。所以这条规则主动放弃优化,宁可留着两个 Filter。
这就是第 09 章"为什么我的谓词没有下推"要展开的东西——大量优化规则都带守卫,守卫不满足时规则静默跳过,你在执行计划里看不到任何提示。
4.5 _ 的多种含义(高频困惑点)
_ 在 Scala 里是个"万能占位符",含义完全取决于上下文。Spark 源码里到处都是,列出你会遇到的几种:
| 写法 | 含义 | Java 等价 |
|---|---|---|
case _ => |
通配符,匹配任何东西 | default: |
_.containsPattern(FILTER) |
匿名函数的参数占位 | x -> x.containsPattern(FILTER) |
TreeNode[_] |
类型通配符,"任意类型参数" | TreeNode<?> |
preCBORules: _* |
把一个集合"摊开"成可变参数 | list.toArray(new X[0]) 传给 varargs |
case Filter(_, child) |
我不关心这个字段 | 无对应,Java 得声明变量再忽略 |
第二种最常见,值得单独说。这两行完全等价:
_.containsPattern(FILTER) // 简写
x => x.containsPattern(FILTER) // 完整写法
x -> x.containsPattern(FILTER) // Java Lambda
规则:_ 在这里表示"传进来的那个参数"。每个 _ 按出现顺序对应一个参数,所以 _ + _ 等价于 (a, b) => a + b。
读源码技巧:看到
_.开头,就在心里替换成x => x.,立刻就通顺了。
五、Option:Scala 版的 Optional
5.1 定义与动机
回到这段代码:
override def maxRows: Option[Long] = condition match {
case Literal.FalseLiteral => Some(0L)
case _ => child.maxRows
}
Option[Long] 表示"一个可能存在、也可能不存在的 Long 值"。它有且仅有两种形态:
Some(值)—— 有值None—— 没有值
Java 等价就是 Optional<Long>:
| Scala | Java |
|---|---|
Option[Long] |
Optional<Long> |
Some(0L) |
Optional.of(0L) |
None |
Optional.empty() |
opt.get |
opt.get() |
opt.getOrElse(默认值) |
opt.orElse(默认值) |
opt.isDefined |
opt.isPresent() |
opt.map(f) |
opt.map(f) |
为什么用它而不是 null:maxRows 的语义是"这个节点最多输出多少行"。有时候能算出来(比如 LIMIT 100 就是 100),有时候算不出来(比如全表扫描)。用 null 表示"算不出来"的问题是——类型系统不强制你检查,忘了判空就是 NullPointerException。而 Option 把"可能没有"编码进类型里,编译器会逼你处理没有值的情况。
5.2 在 Spark 里的实际含义
maxRows 返回 Option[Long] 的业务含义:
Some(0)→ "我确定最多返回 0 行"(比如WHERE false)Some(100)→ "我确定最多返回 100 行"(比如LIMIT 100)None→ "我不知道"(比如一个没有 limit 的全表扫描)
注意 None 不是"0 行",而是"未知"。这个区分在优化器里至关重要:知道"最多 0 行"可以删掉整棵子树;"不知道"则什么都不能做。
排障关联:第 10 章讲统计信息时你会看到,Catalyst 里大量信息都是
Option类型——因为统计信息经常是缺失的。而"统计信息缺失"正是 join 策略选错、广播判断失误的头号原因。
5.3 Option 的常用方法
看这一行(Optimizer.scala:2591):
child.maxRows.exists { _ <= limitExpr.eval().asInstanceOf[Int] }
逐段翻译:
child.maxRows // 得到一个 Option[Long]
.exists { ... } // 如果有值,且值满足条件,返回 true;没值返回 false
{ _ <= ... } // 条件:这个值 <= 某个数
limitExpr.eval() // 计算 LIMIT 后面那个数字
.asInstanceOf[Int] // 强制类型转换成 Int
Java 等价:
child.maxRows() // Optional<Long>
.filter(v -> v <= (Integer) limitExpr.eval()) // Optional 的 filter
.isPresent(); // 有没有剩下值
asInstanceOf[T] 就是 Java 的强制类型转换 (T) x。看到 asInstanceOf 就读作"强转"。
这段代码的业务含义(EliminateLimits 规则):如果子节点最多产出 5 行,而你写了 LIMIT 100,那这个 LIMIT 是多余的,可以直接删掉——因为它永远不会真的截断任何数据。
六、集合操作:对照 Java Stream
Scala 的集合方法和 Java Stream 高度对应。这是你的前置知识里最有用的一项。
6.1 对照表
| Scala | Java Stream | 说明 |
|---|---|---|
list.map(f) |
stream().map(f).toList() |
一对一变换 |
list.filter(f) |
stream().filter(f).toList() |
保留满足条件的 |
list.filterNot(f) |
stream().filter(f.negate()) |
保留不满足的(Java 无直接对应) |
list.flatMap(f) |
stream().flatMap(f) |
变换后摊平 |
list.exists(f) |
stream().anyMatch(f) |
是否存在满足的 |
list.forall(f) |
stream().allMatch(f) |
是否全部满足 |
list.find(f) |
stream().filter(f).findFirst() |
找第一个,返回 Option |
list.partition(f) |
无直接对应 | 按条件分成两组,返回元组 |
list.reduceOption(f) |
stream().reduce(f) |
归约,返回 Option |
list.head |
list.get(0) |
第一个元素 |
list.isEmpty |
list.isEmpty() |
空判断 |
注意 Scala 不需要 .stream() 和 .toList() ——集合方法直接可用,返回同类型集合。这让链式调用比 Java 简洁不少。
6.2 实战:拆解一段真实代码
Optimizer.scala:2038-2039:
val (combineCandidates, rest) =
splitConjunctivePredicates(fc).partition(p => p.deterministic && !p.throwable)
逐段拆:
val (combineCandidates, rest) = ...
// ^^^^^^^^^^^^^^^^^^^^^^^^^ 元组解构:把返回的一对值分别赋给两个变量
splitConjunctivePredicates(fc)
// 把 fc 这个条件按 AND 拆开。比如 "a>1 AND b<2 AND c=3" 拆成 [a>1, b<2, c=3]
// 这个方法来自 PredicateHelper trait —— 还记得 2.5 节说的"方法找不到就去看 trait"吗
.partition(p => p.deterministic && !p.throwable)
// partition:把列表按条件分成两组,返回一个二元组 (满足的, 不满足的)
// 条件:p 是确定性的 并且 p 不会抛异常
Java 等价:
List<Expression> preds = splitConjunctivePredicates(fc);
List<Expression> combineCandidates = preds.stream()
.filter(p -> p.deterministic() && !p.throwable())
.toList();
List<Expression> rest = preds.stream()
.filter(p -> !(p.deterministic() && !p.throwable())) // 取反,要写两遍条件
.toList();
Scala 的 partition 一次遍历分两组,Java 要写两个 stream(或用 Collectors.partitioningBy,但更啰嗦)。
元组(Tuple):(a, b) 是 Scala 内置的"一对值"。Java 没有内置元组,通常要自己定义一个类或用 Map.Entry。看到 val (x, y) = 某方法(),就理解成"这个方法返回两个值,分别赋给 x 和 y"。
七、不可变与 copy:树是怎么"修改"的
7.1 一个必须理解的设计前提
Catalyst 的所有计划节点都是不可变的(case class 的字段默认是 val)。
自问:既然不可变,优化器怎么"修改"计划?比如把两个 Filter 合并成一个?
答:不修改,而是造一个新的。这和 Java 里 String 的行为完全一致——s.replace("a", "b") 不改原字符串,返回新字符串。
7.2 copy:只改一个字段,其余照抄
case class 自动生成的 copy 方法(basicLogicalOperators.scala:357-358):
override protected def withNewChildInternal(newChild: LogicalPlan): Filter =
copy(child = newChild)
翻译:造一个新的 Filter,它的 child 字段换成 newChild,其他字段(condition)从当前对象照抄。
Java 等价:
Filter withNewChildInternal(LogicalPlan newChild) {
return new Filter(this.condition, newChild); // 手动把每个字段列出来
}
copy 的价值在字段多时体现。比如 Join 节点有 6 个字段,只想换一个:
join.copy(condition = newCondition) // Scala:一行
new Join(join.left(), join.right(), join.joinType(),
newCondition, join.hint()); // Java:每个字段都得写,漏一个就是 bug
名字 = 值 是具名参数(named argument)。Java 从 Java 8 起也没有这个特性——具名参数让你不必记住参数顺序,这在 6 个参数的构造函数里是重大的可读性提升。
7.3 为什么坚持不可变(代价与收益)
收益:
- 线程安全。Catalyst 的规则可能并行执行,不可变对象天然无竞态。
- 可以安全共享子树。合并两个 Filter 时,
grandChild整棵子树被新节点直接引用,不需要深拷贝——因为没人能改它。 - 调试友好。计划的每个中间状态都被保留,这就是
explain(mode="extended")能打印出"解析后 → 分析后 → 优化后"三个版本的原因。
代价:
- 对象分配开销。每次规则应用都产生新节点。一个复杂查询经过 200+ 规则,会创建大量临时对象,给 GC 造成压力。这是 Driver 端 CPU 消耗的一个真实来源。
- 深树的重建成本。改一个叶子节点,从根到叶的整条路径上的节点都要重建(因为父节点的
child字段变了)。
这个代价是有实测影响的:超大 SQL(几百个 join、几千行)在 Driver 端优化阶段耗时数分钟甚至 OOM,正是这个机制的直接后果。第 08 章会讲 Catalyst 为此做的优化(
transformWithPruning剪枝、TreePattern位图快速跳过)。
八、把一整条规则读通(综合演练)
现在你已经具备全部前置,来完整读一条真实的优化规则。这是本章的验收点。
文件:sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala:2029-2049
object CombineFilters extends Rule[LogicalPlan] with PredicateHelper {
def apply(plan: LogicalPlan): LogicalPlan = plan.transformWithPruning(
_.containsPattern(FILTER), ruleId)(applyLocally)
val applyLocally: PartialFunction[LogicalPlan, LogicalPlan] = {
case Filter(fc, nf @ Filter(nc, grandChild)) if nc.deterministic =>
val (combineCandidates, rest) =
splitConjunctivePredicates(fc).partition(p => p.deterministic && !p.throwable)
val mergedFilter = (ExpressionSet(combineCandidates) --
ExpressionSet(splitConjunctivePredicates(nc))).reduceOption(And) match {
case Some(ac) =>
Filter(And(nc, ac), grandChild)
case None =>
nf
}
rest.reduceOption(And).map(c => Filter(c, mergedFilter)).getOrElse(mergedFilter)
}
}
8.1 object:Scala 的单例
object CombineFilters extends Rule[LogicalPlan] with PredicateHelper {
object 关键字定义一个单例对象——JVM 里全局只有一个实例。
Java 等价是单例模式,或者一个全是 static 方法的工具类:
// Java 等价(饿汉式单例)
public class CombineFilters extends Rule<LogicalPlan> implements PredicateHelper {
public static final CombineFilters INSTANCE = new CombineFilters();
private CombineFilters() {}
// ...
}
为什么优化规则都是 object:规则是无状态的——它只是"输入一棵树,输出一棵树"的纯函数。无状态就不需要多个实例,单例即可。这也让规则可以被安全地并发调用。
速记:看到
object Xxx,读作"工具类 Xxx,全局唯一实例"。另外,
Rule[LogicalPlan]里的[LogicalPlan]是泛型参数(2.2 节),表示"这条规则处理的树类型是LogicalPlan"。它对应 Java 的Rule<LogicalPlan>。
8.2 apply:可以省略的方法名
def apply(plan: LogicalPlan): LogicalPlan = ...
apply 是 Scala 的一个约定俗成的特殊方法名:如果一个对象有 apply 方法,你可以省略方法名直接调用对象。
CombineFilters.apply(plan) // 完整写法
CombineFilters(plan) // 简写,两者等价
这解释了 2.3 节的一个遗留问题——为什么 case class 可以不写 new:
Filter(cond, child) // 实际调用的是自动生成的 Filter.apply(cond, child)
new Filter(cond, child) // 等价
Java 里没有对应特性。看到 SomeObject(参数) 这种"像调用函数一样调用对象"的写法,翻译成 SomeObject.apply(参数) 就对了。
8.3 PartialFunction:可能"不接受"某些输入的函数
val applyLocally: PartialFunction[LogicalPlan, LogicalPlan] = {
case Filter(fc, nf @ Filter(nc, grandChild)) if nc.deterministic =>
...
}
定义:PartialFunction[A, B](偏函数)是"从 A 到 B 的函数,但只对部分输入有定义"。
注意这段代码只有一个 case,没有 case _ => 兜底。这在普通 match 里会导致运行时异常(没匹配上),但对 PartialFunction 是合法的——没匹配上就表示"这个输入我不处理"。
Java 等价(近似):
// Java 没有偏函数,最接近的是返回 Optional 的函数
Function<LogicalPlan, Optional<LogicalPlan>> applyLocally = plan -> {
if (plan instanceof Filter outer && outer.child() instanceof Filter inner
&& inner.condition().deterministic()) {
return Optional.of(/* 改写后的计划 */);
}
return Optional.empty(); // "我不处理这个输入"
};
为什么 Catalyst 用偏函数:一条规则只关心特定形状的节点。CombineFilters 只关心"Filter 套 Filter",遇到 Join、Aggregate 一律不管。偏函数让规则只声明自己关心的情况,框架负责遍历整棵树并跳过不匹配的节点。
8.4 transformWithPruning:遍历并改写
plan.transformWithPruning(_.containsPattern(FILTER), ruleId)(applyLocally)
逐段:
plan.transformWithPruning( // 遍历整棵树并改写
_.containsPattern(FILTER), // 剪枝条件:子树里含 FILTER 才进去看
ruleId // 规则 id,用于"这条规则已经在这棵子树上跑过了"的标记
)(applyLocally) // 真正的改写逻辑(上面那个偏函数)
注意 )( 这个写法——这是 柯里化(currying):一个方法的参数分成多组,写成 f(a, b)(c) 而不是 f(a, b, c)。
你现在只需要知道:)( 就是参数分了两组,读的时候当成一次调用即可。Spark 里柯里化最常见的用途就是"最后一组参数是个函数",这样可以用大括号写得像代码块:
plan.transformWithPruning(cond, ruleId) {
case Filter(...) => ... // 最后一组参数写成块,可读性更好
}
_.containsPattern(FILTER) 的作用(性能优化,且和第 07 节的代价呼应):整棵计划树可能有几百个节点,而这条规则只关心 Filter。每个节点上维护了一个"我这棵子树里有哪些类型的节点"的位图,如果子树里根本没有 Filter,直接整棵跳过,不用递归进去。这就是 7.3 节提到的"Catalyst 为不可变代价做的优化"。
8.5 完整业务逻辑翻译
现在把整条规则用中文讲一遍:
目标:把嵌套的两层 Filter 合并成一层,减少一次数据遍历。
- 匹配:找到"Filter 套 Filter"的结构,且内层条件是确定性的。
- 拆分外层条件:把外层条件按
AND拆成一个个独立谓词,分成两组:combineCandidates:确定性的、不会抛异常的 → 可以安全合并rest:其余的 → 必须留在外层
- 去重:用
ExpressionSet(...) -- ExpressionSet(...)做集合减法,去掉那些内层已经有的条件(--是集合差集,对应 Java 的removeAll)。比如内层已经有a > 1,外层又写了一遍,就没必要重复判断。 - 合并:把剩下的候选条件用
And连起来,塞进内层 Filter。reduceOption(And):把列表[x, y, z]归约成And(And(x, y), z)。返回Option是因为列表可能为空(全被去重掉了)——这时返回None,走case None => nf分支,直接用原来的内层 Filter。
- 收尾:如果
rest非空(有不能合并的条件),在外面再包一层Filter;否则直接返回合并后的结果。rest.reduceOption(And).map(c => Filter(c, mergedFilter)).getOrElse(mergedFilter)- 翻译:
rest归约成一个条件(可能没有)→ 如果有,包一层 Filter → 如果没有,直接返回mergedFilter
对应的 SQL 效果:
-- 优化前(你写的)
SELECT * FROM (
SELECT * FROM orders WHERE dt = '2026-08-01'
) t WHERE amount > 100
-- 优化后(Catalyst 内部等价于)
SELECT * FROM orders WHERE dt = '2026-08-01' AND amount > 100
8.6 断点快照:如果你在这里打断点,会看到什么
你不在本机编译和调试 Spark,所以这一节替你把"断点看到的东西"写下来。后续每一章的关键位置我都会给这样一张表。
假设执行这条 SQL:
SELECT * FROM (
SELECT * FROM orders WHERE dt = '2026-08-01' AND amount > 0
) t WHERE amount > 100 AND rand() < 0.5
在 applyLocally 的 case 匹配成功那一刻(Optimizer.scala:2037),各个变量的值:
| 变量 | 类型 | 值(toString 的样子) |
说明 |
|---|---|---|---|
fc |
Expression |
((amount#3 > 100) AND (rand() < 0.5)) |
外层 Filter 的条件 |
nf |
Filter |
Filter ((dt#1 = 2026-08-01) AND (amount#3 > 0)) |
内层 Filter 整体 |
nc |
Expression |
((dt#1 = 2026-08-01) AND (amount#3 > 0)) |
内层的条件 |
grandChild |
LogicalPlan |
Relation[orders] |
内层 Filter 的子节点。Relation 就是"读某张表"这个动作对应的叶子节点 |
amount#3里的#3是 exprId——Catalyst 给每个属性分配的全局唯一编号,用来区分同名但不同来源的列。第 06 章(分析器)会讲它怎么来的。你现在只需知道:explain输出里那些#数字是列的身份证号,不是行号。
继续往下走,各步骤的中间值:
| 步骤 | 变量 | 值 | 为什么 |
|---|---|---|---|
| 第 2 步:拆分外层条件 | combineCandidates |
[amount#3 > 100] |
amount > 100 是确定性的,可以合并 |
rest |
[rand() < 0.5] |
rand() 非确定性,必须留在外层 |
|
| 第 3 步:集合去重 | 差集结果 | [amount#3 > 100] |
内层条件里没有它,全部保留 |
| 第 4 步:归约合并 | reduceOption(And) |
Some(amount#3 > 100) |
只有一个元素,归约后就是它自己 |
mergedFilter |
Filter (((dt#1 = ...) AND (amount#3 > 0)) AND (amount#3 > 100)) |
走了 case Some(ac) 分支 |
|
| 第 5 步:收尾 | 最终返回 | Filter (rand() < 0.5) 包住 mergedFilter |
rest 非空,外面再包一层 |
最终的计划树形状:
Filter (rand() < 0.5) ← rest 留在外层
└─ Filter ((dt = '2026-08-01' AND amount > 0) AND amount > 100) ← 合并后的内层
└─ Relation[orders]
这就是"两层 Filter 合并"的真实结果——它没有变成一层。 因为 rand() 挡住了。
这正是排障的价值所在:如果你在
explain里看到两个Filter没被合并,第一反应应该是"是不是有非确定性表达式",而不是"Spark 优化器不行"。第 09 章会把这类"优化为什么没生效"的现场收集成一张排查表。
读到这里,你已经完整读懂了一条真实的 Catalyst 优化规则。 后续 17 章的代码,形态和这条高度一致。
九、读懂一个复杂的方法签名
前面所有例子的方法签名都很短。但 Spark 里有大量长签名,尤其是 core 侧。这是你之前明确指出读不下来的那一类,单独拆。
9.1 目标:读懂这个签名
文件:core/src/main/scala/org/apache/spark/SparkContext.scala:2493
这是 Spark 所有 action 的总入口——你写的每一个 collect()、count()、show(),最终都会走到这里:
def runJob[T, U: ClassTag](
rdd: RDD[T],
func: (TaskContext, Iterator[T]) => U,
partitions: Seq[Int],
resultHandler: (Int, U) => Unit): Unit = {
一共 5 个新语法点。逐个拆。
9.2 [T, U]:两个类型参数
def runJob[T, U: ClassTag](...)
// ^^^^^ 两个泛型参数
<T, U> void runJob(...) // Java 等价
和 Java 泛型方法一样,只是括号形状不同(2.2 节)。
T= RDD 里元素的类型。比如RDD[String]的T就是String。U= 每个分区计算结果的类型。
9.3 A => B:函数类型
这是本节的关键。 Scala 里函数本身是一种类型:
func: (TaskContext, Iterator[T]) => U
// ^^^^^^^^^^^^^^^^^^^^^^^^^ 参数类型(两个参数)
// ^^ 返回类型
读法:「一个函数,它接收一个 TaskContext 和一个 Iterator[T],返回一个 U」。
Java 等价是函数式接口:
BiFunction<TaskContext, Iterator<T>, U> func
对照表(这张表很实用):
| Scala 函数类型 | Java 等价 | 参数数 / 返回 |
|---|---|---|
() => U |
Supplier<U> |
无参,有返回 |
A => U |
Function<A, U> |
1 参,有返回 |
(A, B) => U |
BiFunction<A, B, U> |
2 参,有返回 |
A => Unit |
Consumer<A> |
1 参,无返回 |
(A, B) => Unit |
BiConsumer<A, B> |
2 参,无返回 |
A => Boolean |
Predicate<A> |
1 参,返回布尔 |
Scala 的优势:Java 的函数式接口是一堆有名字的接口(超过 2 个参数就没有内置的了,得自己定义)。Scala 的
(A, B, C, D) => E是统一的语法,参数多少都一样写。
9.4 Unit:Scala 的 void
resultHandler: (Int, U) => Unit
// ^^^^ 没有有意义的返回值
Unit 对应 Java 的 void。方法返回 Unit 意味着它是为了副作用而存在的。
def apply(plan: TreeType): TreeType // 返回一棵树 —— 纯计算
resultHandler: (Int, U) => Unit // 返回 Unit —— 靠副作用干活
resultHandler: (Int, U) => Unit 的业务含义:这是一个回调函数,接收「分区编号 Int」和「该分区的结果 U」,把结果收集到某个地方(比如往数组里塞)。因为它靠往外部数组写值来工作,所以不需要返回值。
严格说:
Unit和void有个技术差异——Unit是一个真实的类型,有唯一的值();void在 Java 里不是类型。这个差异对你读源码没有影响,当void读即可。
9.5 U: ClassTag:上下文界定
这是整个签名里最"劝退"的部分,但它的含义非常具体。
def runJob[T, U: ClassTag](...)
// ^^^^^^^^^^^ 上下文界定(context bound)
先说它不是什么:它不是「U 继承自 ClassTag」。Java 的 <U extends Number> 在 Scala 里写作 [U <: Number](用 <:)。冒号在这里是完全不同的意思。
它是什么:[U: ClassTag] 是一个语法糖,编译器会把它展开成一个额外的隐式参数:
// 你看到的(糖)
def runJob[T, U: ClassTag](rdd: RDD[T], ...): Unit
// 编译器实际生成的(去糖后)
def runJob[T, U](rdd: RDD[T], ...)(implicit evidence: ClassTag[U]): Unit
// ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ 自动多出来的参数组
ClassTag 解决什么问题:JVM 有泛型擦除——运行时 List<String> 和 List<Integer> 是同一个类,类型参数信息丢了。这在 Java 里的后果你可能遇到过:不能写 new T[10]。
Scala 的 ClassTag[U] 就是"把 U 的运行时类型信息偷偷带进来"的机制。Spark 需要它,是因为 runJob 内部要创建 Array[U] 来装结果(回看 2560 行那个重载:返回值就是 Array[U])。
Java 里的对应做法是手动传一个 Class<U> 参数:
// Java 等价:必须手动传 Class 对象
<T, U> void runJob(RDD<T> rdd, Class<U> uClass, ...) {
U[] results = (U[]) Array.newInstance(uClass, n); // 靠 Class 对象创建数组
}
runJob(rdd, String.class, ...); // 调用时必须手动传
Scala 的 ClassTag 就是把这个 Class<U> 参数自动化了——编译器在调用点自动填。
对你读源码的实际影响:
看到
[U: ClassTag],心里读作「这里有个泛型 U,编译器会自动帮忙带上它的运行时类型,我不用管」,然后继续往下读。
它有一个你会遇到的实际后果:ClassTag 是"编译期能推断出具体类型"才拿得到的。这就是为什么 Spark 的某些泛型 API 在动态构造类型的场景下会报 No ClassTag available for U 这类编译错误。你不写 Scala,遇不到;但如果在 Spark 的 issue 或 StackOverflow 里看到这类报错,现在你知道它的来源了。
9.6 隐式参数(implicit):编译器自动填的参数
既然 ClassTag 展开成了隐式参数,顺带把这个概念讲清——Spark 里还有别处用它。
def someMethod(a: Int)(implicit ec: ExecutionContext): Unit
// ^^^^^^^^ 这一组参数你可以不传
规则:标了 implicit 的参数组,调用时可以省略。编译器会在当前作用域里找一个类型匹配、且也标了 implicit 的值自动填进去。
implicit val ec: ExecutionContext = ... // 作用域里有一个隐式的 ExecutionContext
someMethod(1) // 没传第二组参数
someMethod(1)(ec) // 编译器实际调用的
Java 里没有对应特性。 最接近的心智模型是 Spring 的依赖注入——你声明需要某个类型,容器自动找一个填给你。区别是 Spring 在运行时靠反射,Scala 在编译期靠类型查找。
读源码时的处理方式:
看到
(implicit ...)参数组,直接跳过不看。 它对理解方法在干什么几乎没有帮助——它是"编译器的杂务",不是业务逻辑。
9.7 换行的参数列表:只是排版
def runJob[T, U: ClassTag](
rdd: RDD[T],
func: (TaskContext, Iterator[T]) => U,
partitions: Seq[Int],
resultHandler: (Int, U) => Unit): Unit = {
参数多的时候换行写,缩进 4 空格,最后一个参数后面紧跟 ): 返回类型 = {。
注意最后那行的信息密度:
resultHandler: (Int, U) => Unit): Unit = {
// ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ 最后一个参数
// ^ 参数列表结束
// ^^^^ 方法返回类型(Unit)
// ^ = 号,方法体开始
这里有个容易看花眼的地方:Unit): Unit 里出现了两个 Unit——前一个是 resultHandler 这个函数参数的返回类型,后一个是 runJob 自己的返回类型。它们没有关系。
9.8 完整翻译
def runJob[T, U: ClassTag](
rdd: RDD[T],
func: (TaskContext, Iterator[T]) => U,
partitions: Seq[Int],
resultHandler: (Int, U) => Unit): Unit
中文:定义一个泛型方法 runJob,有两个类型参数 T(RDD 元素类型)和 U(分区计算结果类型,且需要 U 的运行时类型信息)。它接收四个参数——要计算的 RDD、一个"给定分区上下文和数据迭代器、产出一个结果"的函数、要计算哪些分区的编号列表、以及一个"拿到分区号和结果后做点什么"的回调。方法本身不返回值。
Java 等价:
<T, U> void runJob(
RDD<T> rdd,
Class<U> uClass, // ClassTag 的手工版
BiFunction<TaskContext, Iterator<T>, U> func,
List<Integer> partitions,
BiConsumer<Integer, U> resultHandler);
业务含义(第 17 章会完整展开):这是 Spark 提交一个 Job 的总入口。func 是「在每个分区上跑什么」,partitions 是「跑哪些分区」,resultHandler 是「结果回到 Driver 后怎么收集」。你写的 df.collect() 最终会构造出这三样东西,然后调到这里。
一个提示:这个方法在
SparkContext.scala里有 7 个重载(2493、2524、2543、2560、2572、2584、2598 行)。其中 2493 行这个是唯一的真实现,其余 6 个都是层层简化的便捷版本,最终都调到它。这是 Spark 源码里非常常见的模式——找"真正干活的那个",就找参数最多的那个重载。
十、剩下的高频语法点(速查)
以下是 Spark 源码里还会遇到、但不值得单独展开的语法。当作词典用。
10.1 字符串插值
s"Starting job: $jobId" // $变量 直接插入
s"count = ${list.size}" // ${表达式} 插入表达式结果
"Starting job: " + jobId // Java 等价
String.format("count = %d", list.size()) // 或者
Spark 里还有一个特殊的 log"..."(结构化日志),第 02 章会说明,你只需知道它也是字符串插值的一种。
10.2 方法调用可以省略点和括号
if (className endsWith "$") ... // 中缀写法
if (className.endsWith("$")) ... // 等价的标准写法
看到 a 方法名 b 这种三段式,翻译成 a.方法名(b)。 这在 Spark 源码里不算多见,但 Rule.scala 里就有一处(见下)。
10.3 无参方法可以省略括号
def ruleName: String = ... // 定义时没有 ()
rule.ruleName // 调用时也不写 ()
Scala 约定:没有副作用的方法(纯查询)定义时就不写括号。这也是为什么你看到 child.output、plan.children 这些看起来像字段的东西,其实可能是方法。
10.4 默认返回最后一个表达式
Scala 的方法不写 return,最后一个表达式的值就是返回值:
val ruleName: String = {
val className = getClass.getName
if (className endsWith "$") className.dropRight(1) else className
// ^ 这个 if-else 的结果就是返回值
}
注意 if-else 在 Scala 里是表达式(有值),不是语句。等价于 Java 的三元运算符:
String className = getClass().getName();
return className.endsWith("$") ? className.substring(0, className.length()-1) : className;
10.5 protected / private[spark]
override protected lazy val validConstraints: ExpressionSet = ...
private[spark] def someMethod() = ...
protected:同 Java。private[spark]:Scala 特有——"仅在spark包及其子包内可见"。这是比 Java 的 package-private 更灵活的访问控制(Java 只能限制在同一个包,不含子包)。
Spark 大量使用 private[spark] 来标记"内部 API,外部用户不要用"。看到它就知道:这是实现细节,不保证版本间稳定。
10.6 可以跳过的东西
以下 Scala 特性在 Spark SQL 侧出现极少,遇到再查即可,不必提前学:
- 隐式转换(
implicit def)——注意和 9.6 节讲的隐式参数是两回事。隐式转换主要用在Column的 DSL 里(让$"name"这种写法成立),读优化器用不到 - 高阶类型、类型类
for推导式(for { x <- ... } yield ...)- 自类型标注(
self: X =>)
十一、边界与坑:Scala 会怎么咬你
这一节是排障导向的。 以下都是 Java 程序员读 Scala 时的真实翻车点。
11.1 == 是值比较,不是引用比较
a == b // Scala:调用 equals,值比较
a.equals(b) // Java 等价
a == b // Java 里这是引用比较!含义完全不同
这是最危险的一个差异。 Scala 的 == 相当于 Java 的 equals(且自动处理 null)。要比较引用得用 eq。
在 Catalyst 里的含义:plan1 == plan2 比较的是两棵树的结构是否完全相同(因为 case class 的 equals 是逐字段递归比较的)。这正是优化器判断"这轮规则有没有改变计划"(不动点判定)的依据——第 08 章会用到。
11.2 找不到方法定义?去看 trait
前面提过,这里再强调一次,因为它是读 Spark 源码时最高频的困惑。
CombineFilters 里调用了 splitConjunctivePredicates,但这个方法在 CombineFilters 里找不到。它来自 with PredicateHelper。
排查顺序:
- 当前类/object 里找
extends的父类里找- 所有
with的 trait 里找 ← 最容易漏 - 父类的父类、trait 的父 trait(继承链可能很深)
11.3 Seq 默认是 List,随机访问是 O(n)
val attrs: Seq[Attribute] = ...
attrs(5) // 看起来像数组下标,实际是 O(n) 的链表遍历
Scala 的 Seq 默认实现是不可变链表(List),不是 ArrayList。attrs(5) 这种"下标访问"实际上要遍历 5 个节点。
在 Spark 里的影响:计划节点的 output、children 都是 Seq。在热路径上按下标随机访问会有性能问题——这也是为什么 Spark 某些地方显式用 IndexedSeq 或 Array。
注意:
attrs(5)这个写法本身也是 8.2 节讲的apply的应用——它实际上是attrs.apply(5)。
11.4 类型推断让你看不到类型
val mergedFilter = (ExpressionSet(...) -- ExpressionSet(...)).reduceOption(And) match {
mergedFilter 是什么类型?代码里没写,靠编译器推断。这对写代码的人方便,对读代码的人是负担。
应对方法:
- 看
match的所有分支返回什么(这里两个分支都返回Filter,所以是LogicalPlan) - 在 IDE 里把光标放上去,会显示推断出的类型
- 你不动手的话,就靠上下文推:看它后面被怎么用
11.5 运算符其实是方法
ExpressionSet(a) -- ExpressionSet(b) // -- 是一个方法名
And(nc, ac) // And 是一个 case class
a + b // + 也是方法
Scala 允许方法名用符号。-- 是集合的差集方法(对应 Java 的 removeAll)。
应对:看到不认识的符号方法,去对应类里搜这个符号。常见的:
++集合拼接(Java 的addAll)--集合差集(removeAll):+/+:追加元素到尾部 / 头部_*展开可变参数(见 4.5 节)
十二、速查表(本章核心产出)
建议收藏这张表,后续章节遇到不认识的写法回来查。
| Scala 写法 | Java 等价 | 一句话说明 |
|---|---|---|
x: Int |
int x |
类型写在后面 |
Seq[Attribute] |
List<Attribute> |
[] 是泛型不是数组 |
def f[T, U](...) |
<T, U> f(...) |
泛型方法 |
[U <: Number] |
<U extends Number> |
类型上界(注意是 <: 不是 :) |
[U : ClassTag] |
手动多传一个 Class<U> |
上下文界定,编译器自动带运行时类型 |
(implicit x: T) |
无对应(近似 Spring 注入) | 隐式参数组,读源码时跳过 |
A => B |
Function<A, B> |
函数类型 |
(A, B) => C |
BiFunction<A, B, C> |
双参函数类型 |
A => Unit |
Consumer<A> |
无返回值的函数 |
Unit |
void |
无有意义的返回值 |
case class X(a: A) |
record X(A a) |
自动生成 equals/hashCode/toString/copy |
object X |
单例类 / 全 static 工具类 | 全局唯一实例 |
trait X |
interface + default 方法 |
可以有字段和实现 |
extends A with B with C |
extends A implements B, C |
第一个 extends,其余 with |
val x = 1 |
final 字段 |
不可变,立即求值 |
var x = 1 |
普通字段 | 可变 |
def x = 1 |
方法 | 每次调用都重算 |
lazy val x = 1 |
双重检查锁的延迟初始化 | 首次访问才算,之后缓存 |
override |
@Override |
Scala 里是强制的 |
private[spark] |
无对应 | 包级可见(含子包) |
x match { case ... } |
switch + instanceof |
可解构对象内部字段 |
case Filter(c, ch) |
if (x instanceof Filter f) + 取字段 |
匹配类型并拆出字段 |
case X(...) if cond |
匹配后再加 if |
守卫条件 |
nf @ Filter(...) |
匹配的同时给整体起名 | @ 是绑定 |
case _ => |
default: |
通配符 |
_.foo |
x -> x.foo() |
匿名函数简写 |
list: _* |
展开成 varargs | 集合摊平成可变参数 |
Option[T] |
Optional<T> |
可能有值可能没有 |
Some(x) / None |
Optional.of(x) / empty() |
有值 / 无值 |
opt.getOrElse(d) |
opt.orElse(d) |
取值或默认 |
x.asInstanceOf[T] |
(T) x |
强制类型转换 |
x.isInstanceOf[T] |
x instanceof T |
类型判断 |
obj(args) |
obj.apply(args) |
apply 可省略 |
Filter(a, b) |
new Filter(a, b) |
case class 免 new |
x.copy(f = v) |
new X(..., v, ...) |
只改一个字段其余照抄 |
(a, b) |
无内置对应 | 元组 |
val (x, y) = f() |
无内置对应 | 元组解构 |
== |
equals() |
值比较,不是引用比较 |
eq |
== |
引用比较 |
s"$x" |
字符串拼接 | 字符串插值 |
f(a)(b) |
f(a, b) |
柯里化,参数分组 |
PartialFunction[A, B] |
返回 Optional 的函数 |
只对部分输入有定义 |