Spark SQL 面试专题
概述
Spark SQL 面试覆盖:DataFrame 与 RDD 关系、UDF 家族、Catalyst/AQE、Shuffle 算子、数据倾斜治理。本文按原理、对比、场景设计三类整理高频问答,附易错点清单。
一、原理类问答
Q1:DataFrame 与 RDD 的区别?
| 维度 | RDD | DataFrame |
|---|---|---|
| Schema | 无 | 有(列类型) |
| 优化 | 无 | Catalyst + Tungsten |
| 序列化 | Java/Kryo | 二进制 |
| 编程体验 | 函数式 | SQL/算子 |
DataFrame 底层仍是 RDD,但经过优化器与代码生成,性能通常更优。
Q2:Dataset 和 DataFrame 什么关系?
DataFrame = Dataset[Row]。Dataset 有强类型(Scala/Java),DataFrame 无类型(适合 SQL/Python)。
Q3:Catalyst 优化器分几个阶段?
Parser → Analyzer → Optimizer → SparkPlanner → Codegen| 阶段 | 产物 |
|---|---|
| Parser | Unresolved 逻辑计划 |
| Analyzer | 绑定 Schema 的逻辑计划 |
| Optimizer | 优化后的逻辑计划 |
| SparkPlanner | 物理计划 |
| Codegen | 执行代码 |
Q4:AQE 自适应优化了哪些?
| 能力 | 说明 |
|---|---|
| 合并 Shuffle 分区 | 减少空任务 |
| 切换 Join 策略 | 运行时改广播 |
| 优化倾斜 Join | 拆分倾斜分区 |
| 动态分区裁剪 | Join 侧裁剪 |
Q5:Spark SQL 为什么比 Hive(MR)快?
| 因素 | 说明 |
|---|---|
| Catalyst 优化 | 谓词下推、列裁剪 |
| Tungsten | 二进制内存 + Codegen |
| 内存计算 | 少落盘 |
| 向量化读 | Parquet 列式快读 |
二、UDF 家族问答
Q6:UDF / UDAF / UDTF 区别?
| 类型 | 输入→输出 | 用途 |
|---|---|---|
| UDF | 一行 → 一行 | 单行处理 |
| UDAF | 多行 → 一行 | 聚合 |
| UDTF | 一行 → 多行 | 行转多列/炸裂 |
sql
-- UDF 示例
spark.udf.register("my_udf", (s: String) => s.trim)
SELECT my_udf(name) FROM t;
-- explode(内置 UDTF)
SELECT explode(split(tags, ',')) FROM t;Q7:如何实现自定义聚合 UDAF?
| 方式 | 说明 |
|---|---|
旧版 UserDefinedAggregateFunction | 继承实现 buffer |
新版 Aggregator[T, BUF, OUT] | 强类型,推荐 |
pandas_udf(Python) | 向量化聚合 |
Q8:UDF 性能问题?
| 问题 | 处理 |
|---|---|
| UDF 阻断优化 | 逻辑尽量用内置函数 |
| Python UDF 慢 | 用向量化 pandas_udf |
| 序列化开销 | 用 Scala/Java UDF |
| null 处理 | UDF 内显式处理 null |
三、Shuffle 算子优化问答
Q9:groupByKey 为什么慢?
groupByKey:全部数据 Shuffle 后聚合
reduceByKey:Map 端先预聚合(combine)Shuffle 数据量:reduceByKey 通常小数倍到数十倍。
Q10:哪些算子会产生 Shuffle?
| 算子 | 说明 |
|---|---|
| reduceByKey / groupByKey | 按 key 重分区 |
| join / cogroup | 关联 |
| distinct / repartition | 去重/重分区 |
| sortByKey | 全局排序 |
Q11:如何减少 Shuffle 数据量?
| 手段 | 说明 |
|---|---|
| 预聚合 | reduceByKey、combineByKey |
| 广播 join | 小表广播免大 Shuffle |
| 分区裁剪 | 只读必要分区 |
| 压缩 | shuffle 压缩 |
| 分桶表 | BucketedJoin 免 Shuffle |
Q12:sortByKey 与 repartitionAndSortWithinPartitions?
| 算子 | 行为 |
|---|---|
| sortByKey | 全局排序(两阶段) |
| repartitionAndSortWithinPartitions | 分区内排序,减少一次排序 |
后者性能更优(Shuffle 中完成排序)。
四、数据倾斜问答
Q13:数据倾斜常见表现?
| 表现 | 原因 |
|---|---|
| 个别 Task 极慢 | key 分布不均 |
| 某 Executor OOM | 单 key 数据量过大 |
| Shuffle 阶段卡住 | 倾斜分区拖累 |
Q14:SQL 场景倾斜怎么治?
| 手段 | 说明 |
|---|---|
| AQE 倾斜优化 | 自动拆分倾斜分区 |
| 加盐 | 热点 key 加随机前缀 |
| 两阶段聚合 | 局部聚合 + 全局聚合 |
| 广播小表 | 换 Join 策略 |
| 提高并行度 | 分散压力 |
Q15:Join 倾斜 vs 聚合倾斜?
| 类型 | 特征 | 对策 |
|---|---|---|
| Join 倾斜 | 热点 key 关联数据多 | AQE 拆分区、加盐、广播 |
| 聚合倾斜 | 单 key 聚合值大 | 两阶段聚合、加盐 |
五、场景设计类问答
Q16:大表 join 大表怎么优化?
方案:
1. 过滤裁剪:join 前先过滤双方无关数据
2. 分桶表:建同分区数的桶表,BucketedJoin
3. 按维度拆分:时间/地域切分再 join
4. AQE:开启倾斜优化
5. 数据预处理:先聚合再 joinQ17:ETL 任务慢,如何排查?
排查顺序:
1. Spark UI 看 Stage 耗时
2. 找慢 Task(倾斜/慢节点)
3. 看 Shuffle 量(预聚合/裁剪)
4. 看 GC(内存配置)
5. explain(true) 查计划(意外 Shuffle/广播)Q18:如何设计分区表?
| 要点 | 说明 |
|---|---|
| 分区列 | 常用过滤条件(dt/hour) |
| 分区粒度 | 天/小时,避免小文件 |
| 分区数 | 适中(几百到几千) |
| 数据倾斜 | 热点分区(如大促日)单独处理 |
六、高频易错点
| 易错点 | 正确理解 |
|---|---|
| DataFrame 无类型 | 运行时才报列错误 |
| SQL null 三值逻辑 | null 比较返回 null 非 false |
| UDF 无法下推 | 谓词下推被自定义函数阻断 |
| AQE 需 ≥3.0 | 2.x 无 AQE |
| cache 后计划固定 | 缓存表绕过多轮优化 |
| partitionBy 写文件多 | 分区粒度细 = 小文件多 |
七、速答清单
| 高频题 | 一句话答案 |
|---|---|
| DataFrame vs RDD | DataFrame 有 Schema + Catalyst 优化 |
| Catalyst 阶段 | Parser/Analyzer/Optimizer/Planner/Codegen |
| AQE 干嘛的 | 运行时合并分区、切 Join、治倾斜 |
| groupByKey vs reduceByKey | reduceByKey 有 map 端预聚合 |
| UDF 家族 | 行/聚合/炸裂三类 |
| 倾斜咋办 | 加盐/两阶段/广播/AQE/提并行度 |