Skip to content

Spark 完全入门指南(下篇):实战调优与前沿 ​

上篇讲原理,中篇讲技术栈和对比,这篇讲最实用的东西:怎么在生产环境用好 Spark,怎么调优,怎么排错,以及 Spark 的未来方向。

面向有经验的工程师,内容偏实战、偏硬核、偏生产环境。


目录 ​


一、生产环境部署与最佳实践 ​

1.1 部署模式选择 ​

YARN 模式(生产环境主流) ​

YARN 有两种提交模式:

模式Driver 位置特点适用场景
client提交任务的机器上日志直接看,方便调试;提交机不能关开发调试
clusterYARN 的 ApplicationMaster 里Driver 由 YARN 管理,有高可用;日志在 YARN 上生产环境

⚠️ 生产环境必须用 cluster 模式: client 模式下,提交任务的机器就是 Driver,如果这台机器挂了或者网络断了,整个任务就失败了。 cluster 模式下,Driver 跑在集群里,YARN 会监控和重试,更可靠。

提交命令示例 ​

bash
# 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 arg2

Kubernetes 模式(云原生方向) ​

bash
# 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.jar

K8s 模式的优势:

  • 弹性伸缩(Dynamic Allocation 在 K8s 上更灵活)
  • 容器化,环境一致性
  • 和云原生生态集成(Prometheus、Grafana、Istio...)
  • GPU 调度支持好

1.2 资源配置最佳实践 ​

Executor 资源配置原则 ​

资源建议原因
Executor 内存4~16GB太小不够用,太大 GC 压力大
Executor Cores2~5 核太多会导致内存竞争,HDFS 客户端瓶颈
Executor 数量根据集群规模不是越多越好,太多 Driver 调度不过来
Driver 内存2~8GBcollect 多的话调大,一般不用太大

常见配置参考 ​

集群规模Executor 内存Executor CoresExecutor 数量总核数
小(3~5 节点)8g3每节点 2 个18~30
中(10~20 节点)12g4每节点 2 个80~160
大(50+ 节点)16g5每节点 2~3 个500+

💡 为什么 Executor Cores 不建议超过 5?

  1. HDFS 客户端并发写入有瓶颈,一般 5 个左右就到上限了
  2. 核太多,多个 Task 共享内存,容易 OOM
  3. 核太多,GC 时所有 Task 都要暂停,影响大

动态资源分配(Dynamic Allocation) ​

properties
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 配置管理最佳实践 ​

配置优先级(从高到低) ​

  1. 代码里 spark.conf.set() / SparkConf.set()
  2. spark-submit 命令行 --conf
  3. 配置文件 spark-defaults.conf
  4. Spark 默认值

💡 最佳实践:

  • 通用配置放 spark-defaults.conf(比如序列化、AQE)
  • 任务特定配置在 spark-submit 命令行指定(比如内存、分区数)
  • 尽量不要在代码里写死配置(不灵活,改配置要重新打包)

重要配置清单 ​

properties
# ===== 基础配置 =====
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-logs

1.4 打包与依赖管理 ​

依赖冲突的坑 ​

Spark 本身带了很多依赖(Guava、Jackson、Netty...),你的应用也可能带这些依赖,版本不一致就会冲突。

解决方案:

  1. Shade 插件(推荐):把依赖重命名打包,避免冲突
  2. provided 范围:Spark 自带的依赖设为 provided,不打进 JAR
  3. 排除依赖:排除冲突的传递依赖

Maven Shade 配置示例 ​

xml
<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 的方法:

  1. 用 reduceByKey 代替 groupByKey

    • reduceByKey:Map 端先聚合(Combiner),Shuffle 数据量小
    • groupByKey:Map 端不聚合,所有数据都要 Shuffle
    • 性能差距可能有好几倍
  2. 小表用 Broadcast Join

    • 把小表广播到每个 Executor,不用 Shuffle
    • 阈值:spark.sql.autoBroadcastJoinThreshold,默认 10MB
    • 手动广播:df.join(broadcast(smallDF), "key")
  3. 尽早过滤

    • 能在 map 端过滤的,不要等到 reduce 端
    • 谓词下推(Catalyst 会自动做,但有时候需要手动优化)
  4. 避免不必要的 distinct

    • distinct 会触发 Shuffle
    • 如果业务上不需要去重,就别用
  5. 用 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()
  • 缓存的数据量太大,把执行内存挤没了 → 合理选择存储级别
  • 缓存了但后续只用一次 → 反而多了一次写的开销,不如不缓存

优化三:数据本地化 ​

数据本地化好的话,计算就在数据所在的节点跑,不用网络传输。

提高数据本地化的方法:

  1. Spark 节点和 HDFS 节点共部(最有效)
  2. 增加 HDFS 副本数(默认 3,可适当增加)
  3. 合理设置本地化等待时间(spark.locality.wait)
  4. 避免把数据缓存在少数几个节点上

优化四:避免 Driver 端计算 ​

很多新手会犯的错误:把大量数据 collect 到 Driver 端处理。

scala
// ❌ 错误:把全量数据拉到 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.partitions200Shuffle 后的分区数根据数据量调整,每个分区 128~256MB
spark.shuffle.file.buffer32kMap 端写 Shuffle 文件的缓冲区内存充足可调大到 64k~128k
spark.reducer.maxSizeInFlight48mReduce 端同时拉取的数据量网络好可调大,减少拉取次数
spark.shuffle.io.maxRetries3Shuffle 拉取重试次数网络不稳定可调大
spark.shuffle.io.retryWait5s重试等待时间网络不稳定可调大
spark.shuffle.compresstrueShuffle 数据是否压缩保持 true,压缩比高
spark.shuffle.managersortShuffle 管理器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 压缩,可以换成更高效的:

properties
# Spark 3.0+ 支持 ZSTD,压缩比和速度都不错
spark.shuffle.compression.codec=zstd
压缩算法压缩比速度说明
lz4低快默认,速度优先
lzf中中老版本默认
snappy中快平衡
zstd高中Spark 3.0+ 支持,推荐
gzip高慢压缩比最高,但慢

2.4 序列化调优 ​

序列化影响 Shuffle 性能和缓存内存占用。

Kryo 序列化 ​

properties
spark.serializer=org.apache.spark.serializer.KryoSerializer
spark.kryoserializer.buffer.max=512m

Kryo 比 Java 序列化:

  • 快 10 倍左右
  • 序列化结果小 2~3 倍
  • 但需要注册类(不注册也能用,只是存类名,占空间)

注册 Kryo 类 ​

scala
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 优化方法 ​

  1. 用 Kryo 序列化:减少对象创建
  2. 序列化存储:MEMORY_ONLY_SER,减少堆内对象
  3. 调大执行内存比例:spark.memory.fraction=0.7(但要小心存储内存不够)
  4. 用 G1 GC:
    properties
    spark.executor.extraJavaOptions=-XX:+UseG1GC -XX:MaxGCPauseMillis=200
  5. 优化代码:减少对象创建,用数组代替集合,用基本类型代替包装类
  6. off-heap 内存:
    properties
    spark.memory.offHeap.enabled=true
    spark.memory.offHeap.size=8g
    用堆外内存,完全避开 GC

2.6 数据格式优化 ​

存储格式对性能影响很大。

格式类型压缩比查询性能适用场景
TextFile行存低差原始数据、日志
CSV行存低差数据交换
JSON行存低差半结构化数据
SequenceFile行存中中Hadoop 生态
Avro行存中中数据序列化、Kafka
Parquet列存高好Spark 默认,OLAP 推荐
ORC列存高好Hive 生态,OLAP

🎯 最佳实践:

  • 中间结果和最终结果用 Parquet 或 ORC(列式存储,压缩比高,查询快)
  • 原始数据可以用 TextFile/JSON(便于查看),但处理前转成列式格式
  • 合理设置分区(按常用查询字段分区,比如 dt)
  • 合理设置分桶(Bucket),对 join 字段分桶可以避免 Shuffle

Parquet 最佳实践 ​

scala
// 写 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 数据量,差异很大 → 倾斜

方法二:看数据分布 ​

scala
// 看 key 的分布
df.groupBy("key").count().orderBy(desc("count")).show(20)

如果前几个 key 的数量远大于其他,就是倾斜。

3.3 数据倾斜的解决方案 ​

方案一:过滤异常 key ​

如果倾斜是由少数异常 key 导致的(比如 null、空字符串、特殊值),直接过滤掉。

scala
// 过滤 null 和空字符串
val filtered = df.filter("key is not null and key != ''")

适用:异常 key 没有业务意义,可以丢弃。

方案二:Broadcast Join(小表 join 大表倾斜) ​

如果是 join 导致的倾斜,而且其中一张表比较小,用 Broadcast Join。

scala
import org.apache.spark.sql.functions.broadcast

// 把小表广播,不用 Shuffle,自然不会倾斜
val result = bigDF.join(broadcast(smallDF), "key")

适用:小表能放进内存(一般 < 1GB)。

方案三:加盐打散(两阶段聚合) ​

这是最经典的倾斜解决方案,适用于 groupBy 或 reduceByKey 倾斜。

思路:

  1. 给倾斜的 key 加上随机前缀(加盐),把一个 key 拆成多个 key
  2. 先做局部聚合(加盐后的 key)
  3. 去掉前缀,再做全局聚合
scala
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 分开处理,再合并。

scala
// 找出倾斜的 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。

properties
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 倾斜。

🎯 数据倾斜调优建议:

  1. 先开 AQE,让 Spark 自动处理(最简单)
  2. AQE 搞不定的,手动分析倾斜原因
  3. 小表 join 用 broadcast
  4. 聚合倾斜用加盐两阶段聚合
  5. 少数 key 倾斜用单独处理
  6. 异常 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 序列化并注册类
scala
// ❌ 错误:在外面创建连接,传到算子里面
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 编码

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)

五、Spark 监控与诊断:Spark UI 深度使用 ​

Spark UI 是调优和排错最重要的工具,一定要会用。

5.1 Spark UI 的页面结构 ​

页面路径作用
Jobs/jobs所有 Job 的概览,看哪个 Job 慢
Stages/stages所有 Stage 的详情,看哪个 Stage 慢、有没有倾斜
Storage/storage缓存的数据,看缓存了什么、占了多少内存
Environment/environment配置信息,确认配置是否生效
Executors/executorsExecutor 列表,看资源使用、GC、Shuffle
SQL/sqlSQL/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。

properties
# 开启事件日志
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 常用的诊断命令 ​

bash
# 查看 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 的意义 ​

  1. 客户端轻量化:不用本地装 Spark,IDE 里直接连远程服务
  2. 多语言支持:不再局限于 Scala/Python/Java/R,任何语言都能用
  3. 云原生友好:服务端可以部署在 K8s 上,弹性伸缩
  4. 统一接口:交互式查询、批处理、流处理用同一个连接

🎯 个人判断: 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 加速主要用于:

  1. 数据预处理/特征工程(特别是图像、文本数据)
  2. 分布式深度学习训练(配合 PyTorch/TensorFlow)
  3. 批量推理(用训练好的模型批量预测)

不是所有场景都适合 GPU,传统的数值计算 CPU 可能更划算。

6.5 湖仓一体:Spark 的新战场 ​

湖仓一体(Lakehouse)是当前大数据最火的方向,Spark 是核心计算引擎。

三大表格式对比(中篇提过,这里深入一点) ​

维度Delta LakeIcebergHudi
创始公司DatabricksNetflixUber
核心设计事务日志(_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(批量推理) ​

python
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 模型。

python
from spark.torch.distributor import TorchDistributor

def train_model():
    # PyTorch 训练代码
    ...

distributor = TorchDistributor(num_processes=4, local_mode=False)
distributor.run(train_model)

方式三:Spark + 向量数据库(RAG) ​

python
# 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 的四大组件:

  1. Tracking:实验跟踪(记录参数、指标、模型)
  2. Projects:项目打包(可复现的运行环境)
  3. Models:模型管理(统一模型格式,多框架支持)
  4. 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 / LlamaIndexRAG 应用开发

🎯 核心观点: 大模型不会取代 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 面试常考知识点(按频率排序) ​

必背(几乎必问) ​

  1. Spark 为什么比 MapReduce 快? —— 内存计算、DAG、丰富 API
  2. 宽窄依赖是什么?怎么划分 Stage? —— 从后往前,遇到宽依赖切分
  3. reduceByKey vs groupByKey —— reduceByKey 有 Combiner,性能好
  4. cache vs persist —— cache 是 persist(MEMORY_ONLY) 的简写
  5. Spark 运行架构 —— Driver + Executor,各自的职责
  6. 数据倾斜怎么解决? —— 过滤、broadcast、加盐、单独处理、AQE

高频(经常问) ​

  1. coalesce vs repartition —— coalesce 默认不 shuffle,repartition 一定 shuffle
  2. Spark 的容错机制 —— Lineage 血缘,checkpoint
  3. Shuffle 过程 —— Map 端写,Reduce 端拉取
  4. 性能调优 —— 代码、参数、资源三个层面
  5. Spark vs Flink —— 批处理 vs 流处理,微批 vs 真流
  6. Spark SQL 优化 —— Catalyst、Tungsten、AQE

中高级岗常问 ​

  1. 内存管理机制 —— 统一内存管理,执行和存储动态借用
  2. DAGScheduler 工作流程 —— Job → Stage → TaskSet
  3. TaskScheduler 调度机制 —— 数据本地性、重试、推测执行
  4. Catalyst 优化器的优化规则 —— 谓词下推、列裁剪、常量折叠...
  5. AQE 的原理和功能 —— 动态合并分区、切换 Join 策略、倾斜优化
  6. Spark on K8s 的原理 —— Driver/Executor 都是 Pod
  7. 湖仓一体和表格式 —— Delta/Iceberg/Hudi 对比
  8. Spark Connect 的架构和意义 —— 客户端-服务端,轻量化

8.3 职业发展方向 ​

掌握 Spark 之后,可以往这些方向发展:

方向说明技能要求
大数据开发工程师最常见,做数据平台、ETL、数仓Spark + Hive + Kafka + 调度
数据仓库工程师做数据建模、数仓建设Spark SQL + 数仓理论 + 维度建模
实时计算工程师做实时数仓、实时推荐Flink + Spark Streaming + Kafka
数据架构师设计大数据架构全栈技术 + 架构设计能力
AI 数据工程师大模型数据处理、RAGSpark + 向量数据库 + LLM 基础
云原生数据工程师云上大数据平台K8s + Spark + 云服务
技术专家/架构师深入源码,做性能优化和架构Spark 源码 + 分布式系统

8.4 推荐资源 ​

官方资源 ​

书籍 ​

  • 《Spark 权威指南》 —— Bill Chambers,比较全面
  • 《深入理解 Spark 核心思想与源码分析》 —— 想深入源码看这本
  • 《Spark 快速大数据分析》 —— 入门经典,有点老

实战项目 ​

  • 日志分析平台
  • 用户行为分析系统
  • 实时数据仓库
  • 推荐系统(召回+排序)
  • 大模型预训练数据处理管道

社区 ​

  • Spark 官方邮件列表
  • GitHub Issues / PR
  • 国内:Spark 技术社区、DataFun、InfoQ 大数据频道

写在最后 ​

Spark 是一个成熟但仍在快速进化的技术。

从 2009 年诞生到现在,它已经从"一个更快的 MapReduce"变成了"统一的大数据计算引擎",又在向"云原生的计算服务"和"AI 时代的数据基础设施"演进。

学 Spark,最重要的是:

  1. 理解原理:知道它为什么这么设计,而不只是会调 API
  2. 动手实践:写代码、调参数、排 Bug,在实战中成长
  3. 关注生态:Spark 不是孤立的,和 Hive、Kafka、Flink、湖仓一体、AI 都有关联
  4. 持续学习:技术在变,AQE、Spark Connect、Photon、GPU、大模型...保持好奇心

大数据的世界很广阔,Spark 是你的船。 愿你在数据的海洋里,航行顺利。

—— 2026 年 8 月


文档版本:V1.0(下篇) 最后更新:2026 年 8 月 适用 Spark 版本:3.5.x

上篇:基础与核心原理 中篇:技术栈与生态对比 下篇:实战调优与前沿(本文)

基于 Vite 强力驱动 | 纯静态轻量托管