Spark 完全入门指南(下篇):实战调优与前沿
上篇讲原理,中篇讲技术栈和对比,这篇讲最实用的东西:怎么在生产环境用好 Spark,怎么调优,怎么排错,以及 Spark 的未来方向。
面向有经验的工程师,内容偏实战、偏硬核、偏生产环境。
目录
- 一、生产环境部署与最佳实践
- 二、性能调优圣经:从入门到精通
- 三、数据倾斜:Spark 调优的头号难题
- 四、排错手册:常见问题与解决方案
- 五、Spark 监控与诊断:Spark UI 深度使用
- 六、前沿方向:Spark 3.x 之后的新世界
- 七、大模型时代的 Spark:AI 浪潮下的定位与进化
- 八、学习路径与职业发展建议
一、生产环境部署与最佳实践
1.1 部署模式选择
YARN 模式(生产环境主流)
YARN 有两种提交模式:
| 模式 | Driver 位置 | 特点 | 适用场景 |
|---|---|---|---|
| client | 提交任务的机器上 | 日志直接看,方便调试;提交机不能关 | 开发调试 |
| cluster | YARN 的 ApplicationMaster 里 | Driver 由 YARN 管理,有高可用;日志在 YARN 上 | 生产环境 |
⚠️ 生产环境必须用 cluster 模式: client 模式下,提交任务的机器就是 Driver,如果这台机器挂了或者网络断了,整个任务就失败了。 cluster 模式下,Driver 跑在集群里,YARN 会监控和重试,更可靠。
提交命令示例
# YARN cluster 模式提交
spark-submit \
--master yarn \
--deploy-mode cluster \
--class com.example.MyApp \
--num-executors 20 \
--executor-cores 4 \
--executor-memory 16g \
--driver-memory 4g \
--conf spark.sql.shuffle.partitions=500 \
--conf spark.sql.adaptive.enabled=true \
--files /path/to/hive-site.xml \
my-app.jar \
arg1 arg2Kubernetes 模式(云原生方向)
# K8s 模式提交
spark-submit \
--master k8s://https://k8s-api-server:6443 \
--deploy-mode cluster \
--class com.example.MyApp \
--conf spark.kubernetes.container.image=spark:3.5.1 \
--conf spark.kubernetes.namespace=spark \
--conf spark.executor.instances=20 \
--conf spark.kubernetes.authenticate.driver.serviceAccountName=spark \
local:///opt/spark/my-app.jarK8s 模式的优势:
- 弹性伸缩(Dynamic Allocation 在 K8s 上更灵活)
- 容器化,环境一致性
- 和云原生生态集成(Prometheus、Grafana、Istio...)
- GPU 调度支持好
1.2 资源配置最佳实践
Executor 资源配置原则
| 资源 | 建议 | 原因 |
|---|---|---|
| Executor 内存 | 4~16GB | 太小不够用,太大 GC 压力大 |
| Executor Cores | 2~5 核 | 太多会导致内存竞争,HDFS 客户端瓶颈 |
| Executor 数量 | 根据集群规模 | 不是越多越好,太多 Driver 调度不过来 |
| Driver 内存 | 2~8GB | collect 多的话调大,一般不用太大 |
常见配置参考
| 集群规模 | Executor 内存 | Executor Cores | Executor 数量 | 总核数 |
|---|---|---|---|---|
| 小(3~5 节点) | 8g | 3 | 每节点 2 个 | 18~30 |
| 中(10~20 节点) | 12g | 4 | 每节点 2 个 | 80~160 |
| 大(50+ 节点) | 16g | 5 | 每节点 2~3 个 | 500+ |
💡 为什么 Executor Cores 不建议超过 5?
- HDFS 客户端并发写入有瓶颈,一般 5 个左右就到上限了
- 核太多,多个 Task 共享内存,容易 OOM
- 核太多,GC 时所有 Task 都要暂停,影响大
动态资源分配(Dynamic Allocation)
spark.dynamicAllocation.enabled=true
spark.dynamicAllocation.minExecutors=2
spark.dynamicAllocation.maxExecutors=50
spark.dynamicAllocation.initialExecutors=5
spark.dynamicAllocation.executorIdleTimeout=60s适用场景:
- 多用户共享集群
- 任务负载波动大
- 交互式查询(Spark Thrift Server)
不适用场景:
- 固定的批处理任务(资源需求稳定,动态分配反而 overhead)
- 对延迟敏感的任务(申请 Executor 需要时间)
1.3 配置管理最佳实践
配置优先级(从高到低)
- 代码里
spark.conf.set()/SparkConf.set() - spark-submit 命令行
--conf - 配置文件
spark-defaults.conf - Spark 默认值
💡 最佳实践:
- 通用配置放
spark-defaults.conf(比如序列化、AQE)- 任务特定配置在 spark-submit 命令行指定(比如内存、分区数)
- 尽量不要在代码里写死配置(不灵活,改配置要重新打包)
重要配置清单
# ===== 基础配置 =====
spark.master=yarn
spark.submit.deployMode=cluster
spark.serializer=org.apache.spark.serializer.KryoSerializer
# ===== 内存配置 =====
spark.executor.memory=12g
spark.executor.cores=4
spark.executor.instances=20
spark.driver.memory=4g
spark.memory.fraction=0.6
spark.memory.storageFraction=0.5
# ===== Spark SQL 配置 =====
spark.sql.shuffle.partitions=500
spark.sql.adaptive.enabled=true
spark.sql.adaptive.coalescePartitions.enabled=true
spark.sql.adaptive.skewJoin.enabled=true
spark.sql.adaptive.localShuffleReader.enabled=true
# ===== 动态分配 =====
spark.dynamicAllocation.enabled=true
spark.dynamicAllocation.minExecutors=2
spark.dynamicAllocation.maxExecutors=50
# ===== 容错配置 =====
spark.task.maxFailures=4
spark.stage.maxConsecutiveAttempts=4
spark.yarn.maxAppAttempts=2
# ===== UI 配置 =====
spark.ui.enabled=true
spark.ui.port=4040
spark.eventLog.enabled=true
spark.eventLog.dir=hdfs:///spark-logs1.4 打包与依赖管理
依赖冲突的坑
Spark 本身带了很多依赖(Guava、Jackson、Netty...),你的应用也可能带这些依赖,版本不一致就会冲突。
解决方案:
- Shade 插件(推荐):把依赖重命名打包,避免冲突
- provided 范围:Spark 自带的依赖设为 provided,不打进 JAR
- 排除依赖:排除冲突的传递依赖
Maven Shade 配置示例
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.5.1</version>
<executions>
<execution>
<phase>package</phase>
<goals><goal>shade</goal></goals>
<configuration>
<filters>
<filter>
<artifact>*:*</artifact>
<excludes>
<exclude>META-INF/*.SF</exclude>
<exclude>META-INF/*.DSA</exclude>
<exclude>META-INF/*.RSA</exclude>
</excludes>
</filter>
</filters>
<relocations>
<relocation>
<pattern>com.google.common</pattern>
<shadedPattern>shaded.com.google.common</shadedPattern>
</relocation>
</relocations>
</configuration>
</execution>
</executions>
</plugin>⚠️ 常见坑:
- 不要把 Spark 本身的依赖打进 fat JAR(用 provided)
- 注意 META-INF 里的签名文件要排除,不然会报 SecurityException
- 服务文件(META-INF/services)要合并,用 Shade 的 ServicesResourceTransformer
二、性能调优圣经:从入门到精通
2.1 调优方法论:先定位,再优化
很多人一上来就调参数,这是错的。正确的调优流程:
1. 看 Spark UI → 找到慢的 Stage / Task
│
▼
2. 分析瓶颈类型 → 计算密集?IO 密集?Shuffle?数据倾斜?
│
▼
3. 针对性优化 → 代码优化 / 参数调优 / 资源调优
│
▼
4. 验证效果 → 对比优化前后的指标
│
▼
5. 重复 → 直到满足性能要求🎯 调优心法: 调优不是瞎调参数,是先找到瓶颈,再针对性优化。 80% 的性能问题是代码和数据的问题,不是参数的问题。 先优化代码和数据,再调参数,最后加资源。
2.2 代码优化(性价比最高)
优化一:减少 Shuffle
Shuffle 是 Spark 性能的第一杀手。能不 Shuffle 就不 Shuffle。
减少 Shuffle 的方法:
用 reduceByKey 代替 groupByKey
- reduceByKey:Map 端先聚合(Combiner),Shuffle 数据量小
- groupByKey:Map 端不聚合,所有数据都要 Shuffle
- 性能差距可能有好几倍
小表用 Broadcast Join
- 把小表广播到每个 Executor,不用 Shuffle
- 阈值:
spark.sql.autoBroadcastJoinThreshold,默认 10MB - 手动广播:
df.join(broadcast(smallDF), "key")
尽早过滤
- 能在 map 端过滤的,不要等到 reduce 端
- 谓词下推(Catalyst 会自动做,但有时候需要手动优化)
避免不必要的 distinct
- distinct 会触发 Shuffle
- 如果业务上不需要去重,就别用
用 coalesce 代替 repartition(减少分区时)
- coalesce 默认不 Shuffle,只是合并分区
- repartition 一定会 Shuffle
- 减少分区用 coalesce,增加分区用 repartition
优化二:用好缓存
什么时候缓存:
- 同一个 DataFrame/RDD 被多次使用
- 计算这个 DataFrame 很耗时
- 数据量适中,放得下内存
缓存级别选择:
| 级别 | 适用场景 |
|---|---|
| MEMORY_ONLY | 数据小,计算快,追求速度 |
| MEMORY_AND_DISK | 数据较大,内存可能放不下(推荐) |
| MEMORY_ONLY_SER | 数据较大,序列化后省内存 |
| MEMORY_AND_DISK_SER | 数据很大,序列化 + 溢写磁盘 |
| DISK_ONLY | 数据非常大,内存完全放不下 |
⚠️ 缓存的坑:
- 缓存了不用的 DataFrame,占着内存不释放 → 记得 unpersist()
- 缓存的数据量太大,把执行内存挤没了 → 合理选择存储级别
- 缓存了但后续只用一次 → 反而多了一次写的开销,不如不缓存
优化三:数据本地化
数据本地化好的话,计算就在数据所在的节点跑,不用网络传输。
提高数据本地化的方法:
- Spark 节点和 HDFS 节点共部(最有效)
- 增加 HDFS 副本数(默认 3,可适当增加)
- 合理设置本地化等待时间(
spark.locality.wait) - 避免把数据缓存在少数几个节点上
优化四:避免 Driver 端计算
很多新手会犯的错误:把大量数据 collect 到 Driver 端处理。
// ❌ 错误:把全量数据拉到 Driver
val allData = df.collect()
val result = allData.map(...).groupBy(...)
// ✅ 正确:用 Spark 的分布式算子
val result = df.groupBy(...).agg(...)⚠️ 记住: Driver 是单点,处理能力有限。 大数据量的计算一定要在 Executor 端分布式完成。 collect() 只用于小数据量的结果查看。
2.3 Shuffle 调优
Shuffle 是最常见的性能瓶颈,专门讲一下。
Shuffle 相关参数
| 参数 | 默认值 | 说明 | 调优建议 |
|---|---|---|---|
spark.sql.shuffle.partitions | 200 | Shuffle 后的分区数 | 根据数据量调整,每个分区 128~256MB |
spark.shuffle.file.buffer | 32k | Map 端写 Shuffle 文件的缓冲区 | 内存充足可调大到 64k~128k |
spark.reducer.maxSizeInFlight | 48m | Reduce 端同时拉取的数据量 | 网络好可调大,减少拉取次数 |
spark.shuffle.io.maxRetries | 3 | Shuffle 拉取重试次数 | 网络不稳定可调大 |
spark.shuffle.io.retryWait | 5s | 重试等待时间 | 网络不稳定可调大 |
spark.shuffle.compress | true | Shuffle 数据是否压缩 | 保持 true,压缩比高 |
spark.shuffle.manager | sort | Shuffle 管理器 | Spark 2.0+ 只有 sort,不用改 |
Shuffle 分区数怎么设?
原则:每个分区大约 128~256MB。
公式:
分区数 ≈ Shuffle 数据总量 / 200MB举例:
- Shuffle 数据量 50GB → 分区数 ≈ 50 * 1024 / 200 ≈ 256 → 设 300
- Shuffle 数据量 500GB → 分区数 ≈ 500 * 1024 / 200 ≈ 2560 → 设 2000~3000
💡 开了 AQE 之后: 分区数可以设大一点,AQE 会自动把小分区合并。 设小了反而可能导致每个分区数据太大。
Shuffle 压缩
默认用 LZF 压缩,可以换成更高效的:
# Spark 3.0+ 支持 ZSTD,压缩比和速度都不错
spark.shuffle.compression.codec=zstd| 压缩算法 | 压缩比 | 速度 | 说明 |
|---|---|---|---|
| lz4 | 低 | 快 | 默认,速度优先 |
| lzf | 中 | 中 | 老版本默认 |
| snappy | 中 | 快 | 平衡 |
| zstd | 高 | 中 | Spark 3.0+ 支持,推荐 |
| gzip | 高 | 慢 | 压缩比最高,但慢 |
2.4 序列化调优
序列化影响 Shuffle 性能和缓存内存占用。
Kryo 序列化
spark.serializer=org.apache.spark.serializer.KryoSerializer
spark.kryoserializer.buffer.max=512mKryo 比 Java 序列化:
- 快 10 倍左右
- 序列化结果小 2~3 倍
- 但需要注册类(不注册也能用,只是存类名,占空间)
注册 Kryo 类
val conf = new SparkConf()
.registerKryoClasses(Array(
classOf[MyClass1],
classOf[MyClass2]
))💡 建议: 生产环境一律用 Kryo 序列化。 自定义类尽量注册,性能更好。 缓冲区大小要够,不然会报 Buffer overflow。
2.5 GC 调优
GC 严重会导致 Task 卡顿,影响性能。
怎么判断 GC 严重?
Spark UI → Task 详情 → GC Time 占比超过 10% 就算严重。
GC 优化方法
- 用 Kryo 序列化:减少对象创建
- 序列化存储:MEMORY_ONLY_SER,减少堆内对象
- 调大执行内存比例:
spark.memory.fraction=0.7(但要小心存储内存不够) - 用 G1 GC:properties
spark.executor.extraJavaOptions=-XX:+UseG1GC -XX:MaxGCPauseMillis=200 - 优化代码:减少对象创建,用数组代替集合,用基本类型代替包装类
- off-heap 内存:properties用堆外内存,完全避开 GC
spark.memory.offHeap.enabled=true spark.memory.offHeap.size=8g
2.6 数据格式优化
存储格式对性能影响很大。
| 格式 | 类型 | 压缩比 | 查询性能 | 适用场景 |
|---|---|---|---|---|
| TextFile | 行存 | 低 | 差 | 原始数据、日志 |
| CSV | 行存 | 低 | 差 | 数据交换 |
| JSON | 行存 | 低 | 差 | 半结构化数据 |
| SequenceFile | 行存 | 中 | 中 | Hadoop 生态 |
| Avro | 行存 | 中 | 中 | 数据序列化、Kafka |
| Parquet | 列存 | 高 | 好 | Spark 默认,OLAP 推荐 |
| ORC | 列存 | 高 | 好 | Hive 生态,OLAP |
🎯 最佳实践:
- 中间结果和最终结果用 Parquet 或 ORC(列式存储,压缩比高,查询快)
- 原始数据可以用 TextFile/JSON(便于查看),但处理前转成列式格式
- 合理设置分区(按常用查询字段分区,比如 dt)
- 合理设置分桶(Bucket),对 join 字段分桶可以避免 Shuffle
Parquet 最佳实践
// 写 Parquet,按 dt 分区
df.write
.partitionBy("dt")
.option("compression", "snappy")
.mode("overwrite")
.parquet("/path/to/data")
// 读 Parquet(自动谓词下推、列裁剪)
val df = spark.read.parquet("/path/to/data")
.filter("dt = '2026-08-11'") // 分区裁剪
.select("col1", "col2") // 列裁剪三、数据倾斜:Spark 调优的头号难题
3.1 什么是数据倾斜?
数据倾斜就是:数据分布不均匀,有的分区数据特别多,有的特别少。
结果就是:
- 大部分 Task 很快跑完
- 少数几个 Task 跑得特别慢(甚至 OOM)
- 整个 Stage 被最慢的 Task 拖死
正常情况:
Task1: ████ (4s)
Task2: ████ (4s)
Task3: ████ (4s)
Task4: ████ (4s)
总时间: 4s
数据倾斜:
Task1: ████ (4s)
Task2: ████████████████████ (20s) ← 倾斜的 Task
Task3: ████ (4s)
Task4: ████ (4s)
总时间: 20s(被最慢的拖死)3.2 怎么发现数据倾斜?
方法一:看 Spark UI
- Stage 详情页,看 Task 的执行时间分布
- 有的 Task 几秒,有的几分钟甚至几小时 → 倾斜
- 看 Shuffle Read 数据量,差异很大 → 倾斜
方法二:看数据分布
// 看 key 的分布
df.groupBy("key").count().orderBy(desc("count")).show(20)如果前几个 key 的数量远大于其他,就是倾斜。
3.3 数据倾斜的解决方案
方案一:过滤异常 key
如果倾斜是由少数异常 key 导致的(比如 null、空字符串、特殊值),直接过滤掉。
// 过滤 null 和空字符串
val filtered = df.filter("key is not null and key != ''")适用:异常 key 没有业务意义,可以丢弃。
方案二:Broadcast Join(小表 join 大表倾斜)
如果是 join 导致的倾斜,而且其中一张表比较小,用 Broadcast Join。
import org.apache.spark.sql.functions.broadcast
// 把小表广播,不用 Shuffle,自然不会倾斜
val result = bigDF.join(broadcast(smallDF), "key")适用:小表能放进内存(一般 < 1GB)。
方案三:加盐打散(两阶段聚合)
这是最经典的倾斜解决方案,适用于 groupBy 或 reduceByKey 倾斜。
思路:
- 给倾斜的 key 加上随机前缀(加盐),把一个 key 拆成多个 key
- 先做局部聚合(加盐后的 key)
- 去掉前缀,再做全局聚合
import org.apache.spark.sql.functions._
// 假设 key 字段倾斜
val salted = df.withColumn("salted_key",
concat(col("key"), lit("_"), (rand() * 10).cast("int")))
// 第一阶段:加盐后局部聚合
val partialAgg = salted.groupBy("salted_key").agg(sum("value") as "partial_sum")
// 第二阶段:去掉盐,全局聚合
val result = partialAgg
.withColumn("key", split(col("salted_key"), "_").getItem(0))
.groupBy("key")
.agg(sum("partial_sum") as "total")适用:聚合操作倾斜,key 数量不多但分布不均。
方案四:倾斜 key 单独处理
如果只有少数几个 key 倾斜,可以把倾斜 key 和正常 key 分开处理,再合并。
// 找出倾斜的 key(比如 count > 100000 的)
val skewKeys = df.groupBy("key").count().filter("count > 100000").select("key")
// 分成倾斜数据和正常数据
val skewDF = df.join(skewKeys, "key")
val normalDF = df.join(skewKeys, Seq("key"), "left_anti")
// 倾斜数据加盐处理
val skewResult = skewDF.withColumn("salt", (rand() * 10).cast("int"))
.groupBy("key", "salt").agg(sum("value") as "v")
.groupBy("key").agg(sum("v") as "total")
// 正常数据正常处理
val normalResult = normalDF.groupBy("key").agg(sum("value") as "total")
// 合并结果
val result = skewResult.union(normalResult)适用:只有少数 key 倾斜,且能识别出来。
方案五:AQE 自动优化(Spark 3.x)
开了 AQE 之后,Spark 会自动检测和优化倾斜 Join。
spark.sql.adaptive.enabled=true
spark.sql.adaptive.skewJoin.enabled=true
spark.sql.adaptive.skewJoin.skewedPartitionFactor=5
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes=256MB原理:AQE 在运行时发现倾斜的分区,会自动把它拆分成多个子分区,分别处理。
适用:Spark 3.x,Join 倾斜。
🎯 数据倾斜调优建议:
- 先开 AQE,让 Spark 自动处理(最简单)
- AQE 搞不定的,手动分析倾斜原因
- 小表 join 用 broadcast
- 聚合倾斜用加盐两阶段聚合
- 少数 key 倾斜用单独处理
- 异常 key 直接过滤
四、排错手册:常见问题与解决方案
4.1 内存相关错误
OOM: Java heap space
Driver OOM:
- 原因:collect() 太多数据、Driver 内存太小、广播变量太大
- 解决:
- 不要 collect 全量数据,用 take 或保存到文件
- 调大
spark.driver.memory - 减少广播变量大小,或者用更高效的序列化
Executor OOM:
- 原因:数据倾斜、单个分区太大、执行内存不够、用户代码创建大对象
- 解决:
- 增加分区数(
repartition或调大 shuffle.partitions) - 处理数据倾斜
- 调大
spark.executor.memory - 用 persist(MEMORY_AND_DISK),内存不够溢写磁盘
- 优化用户代码,不要在 UDF 里创建大对象
- 用 off-heap 内存
- 增加分区数(
OOM: Metaspace
- 原因:JVM 元空间不够,一般是加载的类太多
- 解决:调大
spark.driver.extraJavaOptions=-XX:MaxMetaspaceSize=512m和 Executor 的对应配置
Container killed by YARN for exceeding memory limits
- 原因:Executor 实际使用内存超过了 YARN 分配的 Container 内存
- 解决:
- 调大
spark.executor.memory - 调大
spark.yarn.executor.memoryOverhead(默认是 executor.memory 的 10%,最少 384MB) - 减少每个 Executor 的核数(减少并发 Task)
- 调大
4.2 Shuffle 相关错误
FetchFailedException
- 原因:拉取 Shuffle 数据失败,一般是 Executor 挂了或者网络问题
- 解决:
- 检查 Executor 日志,看是不是 OOM 挂了
- 调大 Shuffle 重试次数和等待时间
- 开启 External Shuffle Service(Executor 挂了 Shuffle 数据还在)
- 处理数据倾斜
MetadataFetchFailedException
- 原因:Shuffle 元数据拉取失败,一般是 Map 端 Executor 挂了
- 解决:同上,主要是保证 Executor 稳定
Too many open files
- 原因:Shuffle 写的文件太多,超过了 Linux 文件句柄限制
- 解决:
- 调大 Linux 文件句柄限制(ulimit -n)
- 用 SortShuffleManager(默认,文件数少)
- 开 AQE 合并小分区,减少 Task 数
4.3 序列化相关错误
NotSerializableException
- 原因:在算子里面用了不能序列化的对象(比如数据库连接、文件句柄)
- 解决:
- 在算子内部创建对象(不要从外面传)
- 用 foreachPartition,在分区级别创建对象
- 让类实现 Serializable
- 用 Kryo 序列化并注册类
// ❌ 错误:在外面创建连接,传到算子里面
val conn = createConnection()
rdd.foreach(x => conn.write(x))
// ✅ 正确:在 foreachPartition 里创建连接
rdd.foreachPartition { iter =>
val conn = createConnection()
iter.foreach(x => conn.write(x))
conn.close()
}Kryo Buffer Overflow
- 原因:Kryo 序列化缓冲区不够大
- 解决:调大
spark.kryoserializer.buffer.max(比如 512m 或 1g)
4.4 数据相关错误
NumberFormatException
- 原因:字符串转数字失败,数据里有非法值
- 解决:
- 数据清洗时处理非法值(过滤或设默认值)
- 用
try-catch或者 Spark 的try_cast(Spark 3.0+)
Schema 不一致
- 原因:读多个文件时,Schema 不一样
- 解决:
- 统一 Schema,读的时候指定 Schema
- 用
mergeSchema选项(Parquet/ORC) - 数据写入时保证 Schema 一致
中文乱码
- 原因:文件编码和读取编码不一致
- 解决:
- 读 CSV 时指定编码:
.option("encoding", "GBK") - 写文件时指定编码
- 统一用 UTF-8 编码
- 读 CSV 时指定编码:
4.5 其他常见错误
ClassNotFoundException / NoClassDefFoundError
- 原因:依赖没打进 JAR,或者版本冲突
- 解决:
- 用 Shade 插件打 fat JAR
- 检查依赖范围(provided 的不会打进去)
- 用
--jars或--packages额外指定依赖 - 检查依赖冲突,用
mvn dependency:tree排查
Task not serializable
- 原因:在算子里面引用了外部不可序列化的对象/方法
- 解决:
- 把用到的变量变成局部变量
- 把方法变成静态方法或者放到可序列化的类里
- 用
@transient标记不需要序列化的字段
ApplicationMaster 挂了 / Driver 挂了
- 原因:Driver OOM、代码异常、资源不足
- 解决:
- 看 Driver 日志(YARN 上用
yarn logs -applicationId <id>) - 调大 Driver 内存
- 检查代码有没有未捕获的异常
- 增加重试次数(
spark.yarn.maxAppAttempts)
- 看 Driver 日志(YARN 上用
五、Spark 监控与诊断:Spark UI 深度使用
Spark UI 是调优和排错最重要的工具,一定要会用。
5.1 Spark UI 的页面结构
| 页面 | 路径 | 作用 |
|---|---|---|
| Jobs | /jobs | 所有 Job 的概览,看哪个 Job 慢 |
| Stages | /stages | 所有 Stage 的详情,看哪个 Stage 慢、有没有倾斜 |
| Storage | /storage | 缓存的数据,看缓存了什么、占了多少内存 |
| Environment | /environment | 配置信息,确认配置是否生效 |
| Executors | /executors | Executor 列表,看资源使用、GC、Shuffle |
| SQL | /sql | SQL/DataFrame 的执行计划和耗时(最常用) |
5.2 关键指标解读
Jobs 页面
- Duration:Job 总耗时
- Stages:Job 包含的 Stage 数
- 点击进去看 Stage 详情
Stages 页面(调优核心)
- Tasks:成功/失败/运行中的 Task 数
- Shuffle Read/Write:Shuffle 数据量
- Duration:Stage 耗时
- 点击进去看 Task 详情:
- Duration:每个 Task 的耗时(看是否倾斜)
- GC Time:GC 时间占比(>10% 就有问题)
- Shuffle Read Size/Records:每个 Task 读了多少数据(看是否倾斜)
- Locality Level:数据本地性级别(PROCESS_LOCAL 最好)
SQL 页面(Spark SQL 调优核心)
- 看执行计划(DAG 可视化)
- 看每个算子的耗时、数据量
- 看有没有 Broadcast Join、有没有 AQE 优化
- 看数据倾斜(某个分区数据量特别大)
Executors 页面
- Storage Memory:存储内存使用
- Executors:每个 Executor 的状态
- Failed Tasks:失败的 Task 数
- GC Time:GC 总时间
- Shuffle Read/Write:Shuffle 总量
5.3 历史服务器(History Server)
任务跑完后 Spark UI 就没了,要查看历史任务需要 History Server。
# 开启事件日志
spark.eventLog.enabled=true
spark.eventLog.dir=hdfs:///spark-logs
spark.eventLog.compress=true
# 启动 History Server
start-history-server.sh访问 http://<server>:18080 查看历史任务。
5.4 常用的诊断命令
# 查看 YARN 应用列表
yarn application -list
# 查看应用日志
yarn logs -applicationId application_1234567890_0001
# 查看应用状态
yarn application -status application_1234567890_0001
# Kill 应用
yarn application -kill application_1234567890_0001
# 查看 HDFS 空间
hdfs dfs -df -h
# 查看 HDFS 文件大小
hdfs dfs -du -h /path/to/data六、前沿方向:Spark 3.x 之后的新世界
6.1 Spark 3.x 重大特性回顾
Spark 3.0(2020)
- Adaptive Query Execution(AQE)—— 运行时动态优化
- Dynamic Partition Pruning —— 动态分区裁剪
- Pandas API on Spark(原 Koalas)—— Pandas 语法分布式执行
- ANSI SQL 兼容模式
- DataSource V2
- Kubernetes 调度 GA
- 30+ 性能优化
Spark 3.1(2021)
- AQE 改进(倾斜 Join 优化)
- Structured Streaming 改进
- Python UDF 性能提升
- DataSource V2 完善
Spark 3.2(2021)
- Pandas API on Spark 集成到 Spark(不再是独立项目)
- AQE 默认开启
- Structured Streaming 状态管理改进
- ANSI SQL 模式完善
Spark 3.3(2022)
- 支持 Java 17
- Pandas API on Spark 完善
- Structured Streaming 触发模式改进(AvailableNow)
- 性能优化
Spark 3.4(2023)
- Spark Connect —— 客户端-服务端架构(重大特性)
- 完整的 ANSI SQL 支持
- Pandas API 完善
- Structured Streaming 改进
Spark 3.5(2023)
- Spark Connect 完善
- 分布式 Deep Learning 集成(sparktorch 等)
- 性能优化
- Python 3.11 支持
6.2 Spark Connect:架构级的变革
Spark Connect 是 Spark 3.4 引入的重大特性,改变了 Spark 的使用方式。
传统架构的问题
客户端(你的代码)+ Driver 绑在一起
- 启动慢(要启动 Driver JVM)
- 资源占用大(客户端要装 Spark)
- 语言支持受限(必须有 Spark 客户端库)
- 多客户端不能共享一个 Spark 集群Spark Connect 架构
轻量客户端 ──gRPC──▶ Spark Connect Server(Driver + 集群)
- 客户端很轻,只需要生成执行计划
- 多客户端共享一个服务端
- 支持更多语言(Go、Rust、JS...只要能生成 protobuf)
- IDE 里写代码不用本地装 Spark
- 更像数据库的使用方式Spark Connect 的意义
- 客户端轻量化:不用本地装 Spark,IDE 里直接连远程服务
- 多语言支持:不再局限于 Scala/Python/Java/R,任何语言都能用
- 云原生友好:服务端可以部署在 K8s 上,弹性伸缩
- 统一接口:交互式查询、批处理、流处理用同一个连接
🎯 个人判断: Spark Connect 是 Spark 走向"云原生化"和"服务化"的重要一步。 未来 Spark 可能会越来越像一个"大数据计算服务",而不是一个"需要自己搭的框架"。 就像从"自己装 MySQL"到"用云数据库 RDS"的转变。
6.3 Photon 引擎:向量化执行的未来
Photon 是 Databricks 公司(Spark 创始团队的公司)开发的向量化执行引擎。
为什么需要 Photon?
Spark 的 Tungsten 虽然做了优化,但还是基于 JVM 的,有几个局限:
- JVM 对象开销(即使是 Unsafe,也有一定 overhead)
- JIT 编译的限制(不能充分利用 CPU 特性)
- GC 压力(即使优化了,还是有)
Photon 的特点
- 用 C++ 写的,完全绕过 JVM
- 向量化执行(一次处理一批数据,不是一条一条)
- 直接操作二进制列式数据
- 充分利用 CPU 缓存和 SIMD 指令
- 支持自适应执行
性能提升
官方数据:Photon 比 JVM 执行引擎快 2~10 倍,特别是在 SQL 查询和数据清洗场景。
现状
- 目前是 Databricks 商业版的核心功能
- 开源版 Spark 还没有完全跟进
- 但开源社区也在做向量化执行的探索(比如 Arrow 集成)
💡 趋势判断: 向量化执行是数据库/大数据引擎的大趋势。 ClickHouse、DuckDB、Doris 都是向量化执行,性能非常猛。 Spark 如果不跟进,在 OLAP 场景会被拉开差距。 未来 Photon 会不会开源?这是很多人关注的问题。
6.4 GPU 加速:AI 时代的必然
NVIDIA RAPIDS Accelerator for Spark
NVIDIA 开发的 Spark GPU 加速插件,可以让 Spark 利用 GPU 计算。
原理:
- 把部分算子(过滤、投影、聚合、Join、排序等)下推到 GPU 执行
- 数据用 GPU 内存(显存)
- 利用 GPU 的并行计算能力
性能:
- 某些场景(数据清洗、ETL、机器学习特征工程)能快 2~5 倍
- 成本可能更低(GPU 贵但算得快,总耗时少)
现状:
- 还在发展中,不是所有算子都支持 GPU
- 需要 NVIDIA GPU 和 CUDA 环境
- 配置相对复杂
Spark + Deep Learning
- TorchDistributor:Spark 3.5 引入,分布式 PyTorch 训练
- sparktorch:第三方库,在 Spark 上训练 PyTorch 模型
- Horovod:Uber 开源的分布式深度学习框架,支持 Spark
🎯 GPU 加速的定位: Spark 的 GPU 加速主要用于:
- 数据预处理/特征工程(特别是图像、文本数据)
- 分布式深度学习训练(配合 PyTorch/TensorFlow)
- 批量推理(用训练好的模型批量预测)
不是所有场景都适合 GPU,传统的数值计算 CPU 可能更划算。
6.5 湖仓一体:Spark 的新战场
湖仓一体(Lakehouse)是当前大数据最火的方向,Spark 是核心计算引擎。
三大表格式对比(中篇提过,这里深入一点)
| 维度 | Delta Lake | Iceberg | Hudi |
|---|---|---|---|
| 创始公司 | Databricks | Netflix | Uber |
| 核心设计 | 事务日志(_delta_log) | 元数据树(metadata + manifest) | 时间线 + 文件索引 |
| Upsert 性能 | 好(MERGE INTO) | 好(MERGE INTO) | 最好(原生 Upsert,索引优化) |
| 流式写入 | 好 | 好 | 好 |
| 流式读取 | 好(CDC) | 好 | 好(增量读取) |
| 时间旅行 | ✅ | ✅ | ✅ |
| Schema 演进 | ✅ | ✅(更灵活) | ✅ |
| 多引擎支持 | Spark 最好 | Spark/Flink/Trino 都好 | Spark/Flink 好 |
| 国内大厂使用 | 越来越多 | 腾讯、阿里等 | 阿里、字节、腾讯等 |
| 社区活跃度 | 高 | 高 | 高 |
Spark 在湖仓一体中的角色
- 计算引擎:读写表格式数据、执行 SQL、ETL
- 流处理:Structured Streaming 写入和读取
- 机器学习:特征工程、模型训练数据准备
- 数据质量:数据校验、监控
🎯 湖仓一体的未来: 湖仓一体正在成为大数据架构的新标准。 传统的"数据湖 + 数据仓库"两套架构会逐渐被"湖仓一体"替代。 Spark + 表格式(Delta/Iceberg/Hudi)是目前最主流的湖仓一体方案。 作为大数据工程师,湖仓一体是必须掌握的方向。
七、大模型时代的 Spark:AI 浪潮下的定位与进化
2023 年以来,大模型(LLM)火遍全球。Spark 在这个新时代扮演什么角色?
7.1 Spark 不擅长什么:大模型训练
首先明确:Spark 不是大模型训练框架。
大模型训练需要:
- 分布式深度学习框架(PyTorch DDP/FSDP、DeepSpeed、Megatron)
- GPU 集群(A100/H100 等)
- 模型并行、张量并行、流水线并行
- 混合精度训练、梯度累积等技术
这些都不是 Spark 的强项。Spark MLlib 只支持传统机器学习,不支持深度学习。
💡 别搞错定位: 不要用 Spark 去训练大模型,那是用错工具了。 大模型训练用 PyTorch + DeepSpeed/Megatron。
7.2 Spark 擅长什么:大模型的数据基础设施
大模型训练和应用中,80% 的工作是数据处理。这正是 Spark 的主场。
场景一:预训练数据处理
- 爬取海量文本数据(网页、书籍、代码...)
- 数据清洗(去重、去噪、过滤低质量)
- 数据格式化(转成模型需要的格式)
- 数据混合(不同来源按比例混合)
- Tokenization(分词,生成 token id)
这些数据量通常是 TB~PB 级,单机处理不了,Spark 是主力工具。
场景二:微调数据准备
- SFT(监督微调)数据构造
- RLHF(人类反馈强化学习)数据准备
- 偏好数据对构造
- 数据清洗和质量控制
场景三:RAG(检索增强生成)数据处理
- 文档解析和分块(Chunking)
- 批量生成 Embedding(可以配合 GPU)
- 向量入库(写入向量数据库)
- 知识库更新和维护
场景四:批量推理(Batch Inference)
- 用训练好的大模型批量处理数据
- 比如批量生成摘要、批量翻译、批量分类
- Spark 负责数据读取和分发,模型推理可以在 Executor 上跑(CPU 或 GPU)
场景五:模型评估
- 批量生成测试集的预测结果
- 计算评估指标(准确率、BLEU、ROUGE...)
- 对比不同模型的效果
7.3 Spark + 大模型的集成方式
方式一:Spark + Pandas UDF(批量推理)
import pandas as pd
from transformers import pipeline
@pandas_udf("string")
def summarize(texts: pd.Series) -> pd.Series:
# 在 Executor 上加载模型(每个分区加载一次)
summarizer = pipeline("summarization", model="facebook/bart-large-cnn")
results = summarizer(texts.tolist(), max_length=50)
return pd.Series([r["summary_text"] for r in results])
# 批量推理
df.withColumn("summary", summarize(col("content"))).show()方式二:Spark + TorchDistributor(分布式训练)
Spark 3.5 引入的 TorchDistributor,可以在 Spark 上分布式训练 PyTorch 模型。
from spark.torch.distributor import TorchDistributor
def train_model():
# PyTorch 训练代码
...
distributor = TorchDistributor(num_processes=4, local_mode=False)
distributor.run(train_model)方式三:Spark + 向量数据库(RAG)
# 1. 用 Spark 批量处理文档
documents = spark.read.text("/path/to/docs")
chunks = documents.withColumn("chunks", split_text(col("value")))
# 2. 批量生成 Embedding(用 Pandas UDF + 模型)
embeddings = chunks.withColumn("embedding", generate_embedding(col("chunk")))
# 3. 写入向量数据库(Milvus/Qdrant/Weaviate...)
embeddings.write.format("milvus").save()7.4 MLflow:Spark 团队的 AI 平台作品
MLflow 是 Databricks(Spark 创始团队)开源的 ML 平台,现在是 MLOps 的事实标准。
MLflow 的四大组件:
- Tracking:实验跟踪(记录参数、指标、模型)
- Projects:项目打包(可复现的运行环境)
- Models:模型管理(统一模型格式,多框架支持)
- Registry:模型注册(版本管理、阶段管理)
💡 Spark 和 MLflow 的关系: 都是 Databricks 搞的,天然集成。 Spark 做数据处理和模型训练,MLflow 做实验管理和模型部署。 做 AI 项目,Spark + MLflow 是黄金搭档。
7.5 大模型时代的技能栈建议
作为大数据/数据工程师,在大模型时代需要掌握:
| 层级 | 技能 | 说明 |
|---|---|---|
| 数据层 | Spark / Flink | 数据处理,基本功 |
| 存储层 | 数据湖 + 表格式(Delta/Iceberg/Hudi) | 湖仓一体 |
| 向量层 | 向量数据库(Milvus/Qdrant) | RAG 必备 |
| 模型层 | PyTorch / Hugging Face | 大模型基础,不用精通但要会用 |
| 编排层 | MLflow / Airflow | 实验管理和工作流 |
| 应用层 | LangChain / LlamaIndex | RAG 应用开发 |
🎯 核心观点: 大模型不会取代 Spark,反而会增加对 Spark 的需求(数据处理量更大了)。 但只会 Spark 不够了,需要往 AI 方向扩展技能。 数据工程师 + AI 工程能力,是未来几年的热门方向。
八、学习路径与职业发展建议
8.1 Spark 学习路径(给有经验的工程师)
第一阶段:快速上手(1~2 周)
- 搭环境,跑通 WordCount
- 学习 DataFrame/SQL API(工作中最常用)
- 理解核心概念:RDD、转换/行动算子、宽窄依赖
- 会看 Spark UI
第二阶段:深入理解(2~4 周)
- 理解 DAG、Stage、Task 的调度机制
- 理解内存管理、Shuffle 过程
- 学习 Structured Streaming
- 学习性能调优基础
- 做 1~2 个实战项目
第三阶段:调优与排错(1~2 个月)
- 深入学习性能调优(代码、参数、资源)
- 学习数据倾斜处理
- 学习常见错误排查
- 学习 Spark 监控和诊断
- 在实际项目中积累经验
第四阶段:源码与进阶(可选,3~6 个月)
- 阅读 Spark 核心源码(DAGScheduler、TaskScheduler、Shuffle、内存管理)
- 学习 Catalyst 优化器
- 学习 Tungsten 执行引擎
- 自定义数据源、自定义优化规则
第五阶段:生态扩展(持续)
- 学习 Flink(流处理方向)
- 学习湖仓一体(Delta/Iceberg/Hudi)
- 学习云原生(Spark on K8s)
- 学习 AI 工程(MLflow、大模型数据处理)
8.2 面试常考知识点(按频率排序)
必背(几乎必问)
- Spark 为什么比 MapReduce 快? —— 内存计算、DAG、丰富 API
- 宽窄依赖是什么?怎么划分 Stage? —— 从后往前,遇到宽依赖切分
- reduceByKey vs groupByKey —— reduceByKey 有 Combiner,性能好
- cache vs persist —— cache 是 persist(MEMORY_ONLY) 的简写
- Spark 运行架构 —— Driver + Executor,各自的职责
- 数据倾斜怎么解决? —— 过滤、broadcast、加盐、单独处理、AQE
高频(经常问)
- coalesce vs repartition —— coalesce 默认不 shuffle,repartition 一定 shuffle
- Spark 的容错机制 —— Lineage 血缘,checkpoint
- Shuffle 过程 —— Map 端写,Reduce 端拉取
- 性能调优 —— 代码、参数、资源三个层面
- Spark vs Flink —— 批处理 vs 流处理,微批 vs 真流
- Spark SQL 优化 —— Catalyst、Tungsten、AQE
中高级岗常问
- 内存管理机制 —— 统一内存管理,执行和存储动态借用
- DAGScheduler 工作流程 —— Job → Stage → TaskSet
- TaskScheduler 调度机制 —— 数据本地性、重试、推测执行
- Catalyst 优化器的优化规则 —— 谓词下推、列裁剪、常量折叠...
- AQE 的原理和功能 —— 动态合并分区、切换 Join 策略、倾斜优化
- Spark on K8s 的原理 —— Driver/Executor 都是 Pod
- 湖仓一体和表格式 —— Delta/Iceberg/Hudi 对比
- Spark Connect 的架构和意义 —— 客户端-服务端,轻量化
8.3 职业发展方向
掌握 Spark 之后,可以往这些方向发展:
| 方向 | 说明 | 技能要求 |
|---|---|---|
| 大数据开发工程师 | 最常见,做数据平台、ETL、数仓 | Spark + Hive + Kafka + 调度 |
| 数据仓库工程师 | 做数据建模、数仓建设 | Spark SQL + 数仓理论 + 维度建模 |
| 实时计算工程师 | 做实时数仓、实时推荐 | Flink + Spark Streaming + Kafka |
| 数据架构师 | 设计大数据架构 | 全栈技术 + 架构设计能力 |
| AI 数据工程师 | 大模型数据处理、RAG | Spark + 向量数据库 + LLM 基础 |
| 云原生数据工程师 | 云上大数据平台 | K8s + Spark + 云服务 |
| 技术专家/架构师 | 深入源码,做性能优化和架构 | Spark 源码 + 分布式系统 |
8.4 推荐资源
官方资源
书籍
- 《Spark 权威指南》 —— Bill Chambers,比较全面
- 《深入理解 Spark 核心思想与源码分析》 —— 想深入源码看这本
- 《Spark 快速大数据分析》 —— 入门经典,有点老
实战项目
- 日志分析平台
- 用户行为分析系统
- 实时数据仓库
- 推荐系统(召回+排序)
- 大模型预训练数据处理管道
社区
- Spark 官方邮件列表
- GitHub Issues / PR
- 国内:Spark 技术社区、DataFun、InfoQ 大数据频道
写在最后
Spark 是一个成熟但仍在快速进化的技术。
从 2009 年诞生到现在,它已经从"一个更快的 MapReduce"变成了"统一的大数据计算引擎",又在向"云原生的计算服务"和"AI 时代的数据基础设施"演进。
学 Spark,最重要的是:
- 理解原理:知道它为什么这么设计,而不只是会调 API
- 动手实践:写代码、调参数、排 Bug,在实战中成长
- 关注生态:Spark 不是孤立的,和 Hive、Kafka、Flink、湖仓一体、AI 都有关联
- 持续学习:技术在变,AQE、Spark Connect、Photon、GPU、大模型...保持好奇心
大数据的世界很广阔,Spark 是你的船。 愿你在数据的海洋里,航行顺利。
—— 2026 年 8 月
文档版本:V1.0(下篇) 最后更新:2026 年 8 月 适用 Spark 版本:3.5.x
上篇:基础与核心原理 中篇:技术栈与生态对比 下篇:实战调优与前沿(本文)