Spark SQL 架构
概述
Spark SQL 是 Spark 上执行 SQL 与结构化数据处理的模块,核心由 Catalyst 优化器与 Tungsten 执行引擎组成。本文讲清 DataFrame/Dataset/RDD 的关系、SQL 到物理执行的全流程,以及 Catalyst/Tungsten 如何带来性能飞跃。
一、Spark SQL 是什么
| 能力 | 说明 |
|---|---|
| SQL 查询 | 标准 SQL 与 HiveQL |
| DataFrame API | 结构化数据编程接口 |
| 数据源统一 | 读 Hive/Parquet/JDBC/JSON 等 |
| 与代码整合 | SQL 与 Scala/Java/Python 混写 |
一句话:让开发者用 SQL 的简洁与代码的灵活,同时享受优化器的性能。
二、RDD / DataFrame / Dataset
2.1 三者对比
| 维度 | RDD | DataFrame | Dataset |
|---|---|---|---|
| 类型安全 | 强类型 | 无(Row) | 强类型 |
| Schema | 无 | 有 | 有 |
| 优化 | 无(用户控制) | Catalyst 优化 | Catalyst 优化 |
| 序列化 | Java/Kryo | 二进制(Tungsten) | 二进制 |
| 编译期检查 | 无 | 无 | 有(Scala/Java) |
2.2 关系
Dataset[T](强类型)
│ 底层是 DataFrame(Dataset[Row])
↓
DataFrame = Dataset[Row]
│ 执行时转为 RDD + Catalyst 优化计划
↓
RDD(最底层执行模型)| 结论 | 说明 |
|---|---|
| DataFrame 是 Dataset[Row] | 无类型,适合 SQL 与 Python |
| Dataset 强类型 | 适合 Scala/Java 复杂逻辑 |
| 三者可互转 | df.rdd / df.as[CaseClass] |
| DataFrame 性能 ≥ RDD | Catalyst/Tungsten 优化 |
三、SQL 执行全流程
3.1 流程总览
SQL / DataFrame API
│
1. SQL 解析(Parser):SQL → Unresolved Logical Plan
│
2. 语义分析(Analyzer):绑定 Schema → Logical Plan
│
3. 逻辑优化(Optimizer):规则优化 → Optimized Logical Plan
│
4. 物理计划(SparkPlanner):选择执行策略 → Physical Plan
│
5. 生成代码(Tungsten):WholeStageCodegen → RDD 执行
│
6. 提交执行:Job → Task 并行计算3.2 各阶段产物
| 阶段 | 产物 | 作用 |
|---|---|---|
| Parser | Unresolved Logical Plan | 语法树,未绑定列 |
| Analyzer | Logical Plan | 列/表绑定 Schema,类型检查 |
| Optimizer | Optimized Logical Plan | 规则优化(谓词下推等) |
| SparkPlanner | Physical Plan | 选择物理算子与策略 |
| Codegen | Java 代码 | 生成高效执行代码 |
四、Catalyst 优化器
4.1 组成
| 组件 | 职责 |
|---|---|
| Parser | SQL 字符串 → 语法树 |
| Analyzer | 绑定 Catalog 与 Schema |
| Optimizer | 逻辑规则优化 |
| SparkPlanner | 逻辑 → 物理计划 |
| 执行准备 | 代码生成、物理算子转换 |
4.2 逻辑优化规则(示例)
| 规则 | 说明 |
|---|---|
| 谓词下推 | WHERE 条件尽量提前到数据源 |
| 列裁剪 | 只读需要的列 |
| 常量折叠 | 常量表达式预先计算 |
| 消除无用操作 | 去重冗余的投影/过滤 |
| 合并算子 | 连续过滤/投影合并 |
4.3 物理优化
| 策略 | 说明 |
|---|---|
| Join 策略选择 | 按表大小选 Broadcast/SortMerge/Shuffle |
| 数据源下推 | 过滤/裁剪下推到 JDBC/Parquet |
| 分区裁剪 | 跳过不相关分区 |
五、Tungsten 执行引擎
5.1 三大优化
| 优化 | 说明 |
|---|---|
| 二进制内存管理 | 数据以二进制存储,避免 Java 对象开销 |
| 缓存友好布局 | 列式/紧凑布局,提升 Cache 命中 |
| WholeStage Codegen | 把整个 Stage 编译成一段 Java 代码 |
5.2 WholeStageCodegen
传统执行:每个算子一个虚函数调用,产生大量对象
Codegen:整个 Stage 的算子融合为一个 Java 方法循环
例:
SELECT sum(x) FROM t WHERE x > 10
生成:
while (iter.hasNext) {
if (row.x > 10) sum += row.x; // 无中间对象
}| 收益 | 说明 |
|---|---|
| 消除虚函数调用 | 提升数倍吞吐 |
| 减少对象分配 | 降低 GC 压力 |
| CPU 友好 | 顺序内存访问 |
六、编程模型
6.1 DataFrame 示例
scala
val df = spark.read.parquet("hdfs:///data/orders")
df.createOrReplaceTempView("orders")
val result = spark.sql("""
SELECT user_id, SUM(amount) AS total
FROM orders
WHERE dt = '2026-08-01'
GROUP BY user_id
ORDER BY total DESC
""")6.2 三种 API 对应
| API | 写法 |
|---|---|
| SQL | spark.sql("SELECT ...") |
| DataFrame | df.select(...).groupBy(...).agg(...) |
| Dataset | df.as[Order].filter(_.dt == "2026-08-01") |
三者最终都走 Catalyst + Tungsten 同一套执行引擎。
七、常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 列不存在报错 | 表 Schema 与列名不匹配,检查大小写 |
| 小表 join 大表慢 | 用 Broadcast 提示或自动优化 |
| 类型不匹配 | Row 与 Dataset 强类型转换注意 |
| SQL 与代码结果不一致 | 检查 null 处理(SQL 三值逻辑) |
| 看不到优化效果 | 用 df.explain(true) 查物理计划 |