全文共 8,874 字 预计阅读 26 分钟
spark

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 有几个现实原因:

  1. Scala 运行在 JVM 上,能直接复用整个 Hadoop / Java 生态(HDFS 客户端、序列化库等)。这是硬约束——大数据生态的既有资产全在 JVM 上。
  2. Scala 有 REPL(交互式解释器)spark-shell 能存在,是因为 Scala 自带 REPL。2010 年前后的 Java 没有(Java 9 才有 jshell),而"交互式探索数据"是 Spark 早期最重要的卖点之一。
  3. 函数式语法让 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 大量用它的原因):

  1. 可以有字段(Java interface 只能有 public static final 常量)
  2. 可以有构造逻辑
  3. 一个类可以混入任意多个 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 类,这里只为表达语义

对应到你写的 SQLWHERE 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 > 11 < 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

overridefinal

final override val nodePatterns: Seq[TreePattern] = Seq(FILTER)
  • overrideScala 里覆写父类成员必须显式写 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";
};

对应关系:

  • matchswitch
  • case 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)

为什么用它而不是 nullmaxRows 的语义是"这个节点最多输出多少行"。有时候能算出来(比如 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 为什么坚持不可变(代价与收益)

收益

  1. 线程安全。Catalyst 的规则可能并行执行,不可变对象天然无竞态。
  2. 可以安全共享子树。合并两个 Filter 时,grandChild 整棵子树被新节点直接引用,不需要深拷贝——因为没人能改它。
  3. 调试友好。计划的每个中间状态都被保留,这就是 explain(mode="extended") 能打印出"解析后 → 分析后 → 优化后"三个版本的原因。

代价

  1. 对象分配开销。每次规则应用都产生新节点。一个复杂查询经过 200+ 规则,会创建大量临时对象,给 GC 造成压力。这是 Driver 端 CPU 消耗的一个真实来源。
  2. 深树的重建成本。改一个叶子节点,从根到叶的整条路径上的节点都要重建(因为父节点的 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",遇到 JoinAggregate 一律不管。偏函数让规则只声明自己关心的情况,框架负责遍历整棵树并跳过不匹配的节点。

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 合并成一层,减少一次数据遍历。

  1. 匹配:找到"Filter 套 Filter"的结构,且内层条件是确定性的。
  2. 拆分外层条件:把外层条件按 AND 拆成一个个独立谓词,分成两组:
    • combineCandidates:确定性的、不会抛异常的 → 可以安全合并
    • rest:其余的 → 必须留在外层
  3. 去重:用 ExpressionSet(...) -- ExpressionSet(...) 做集合减法,去掉那些内层已经有的条件-- 是集合差集,对应 Java 的 removeAll)。比如内层已经有 a > 1,外层又写了一遍,就没必要重复判断。
  4. 合并:把剩下的候选条件用 And 连起来,塞进内层 Filter。
    • reduceOption(And):把列表 [x, y, z] 归约成 And(And(x, y), z)返回 Option 是因为列表可能为空(全被去重掉了)——这时返回 None,走 case None => nf 分支,直接用原来的内层 Filter。
  5. 收尾:如果 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

applyLocallycase 匹配成功那一刻(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 里的 #3exprId——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」,把结果收集到某个地方(比如往数组里塞)。因为它靠往外部数组写值来工作,所以不需要返回值。

严格说Unitvoid 有个技术差异——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.outputplan.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 classequals 是逐字段递归比较的)。这正是优化器判断"这轮规则有没有改变计划"(不动点判定)的依据——第 08 章会用到。

11.2 找不到方法定义?去看 trait

前面提过,这里再强调一次,因为它是读 Spark 源码时最高频的困惑

CombineFilters 里调用了 splitConjunctivePredicates,但这个方法在 CombineFilters 里找不到。它来自 with PredicateHelper

排查顺序

  1. 当前类/object 里找
  2. extends 的父类里找
  3. 所有 with 的 trait 里找 ← 最容易漏
  4. 父类的父类、trait 的父 trait(继承链可能很深)

11.3 Seq 默认是 List,随机访问是 O(n)

val attrs: Seq[Attribute] = ...
attrs(5)          // 看起来像数组下标,实际是 O(n) 的链表遍历

Scala 的 Seq 默认实现是不可变链表List),不是 ArrayListattrs(5) 这种"下标访问"实际上要遍历 5 个节点。

在 Spark 里的影响:计划节点的 outputchildren 都是 Seq。在热路径上按下标随机访问会有性能问题——这也是为什么 Spark 某些地方显式用 IndexedSeqArray

注意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 的函数 只对部分输入有定义
Back to Blog