Spark故障排除与性能优化实战指南
Spark故障排除与性能优化实战指南
一、运维管理
1. Master节点故障与 standby 重启无效
Master 进程默认仅分配 512MB 内存,当集群中运行大量任务时,Master 需要读取每个 Task 的事件日志以生成 SparkUI,内存不足会导致 OOM。通过 HA 方式启动的 Master 同样会因此失败。
解决方案
方案一:增加 Master 内存配置,在 Master 节点配置 spark-env.sh:
export SPARK_DAEMON_MEMORY 8g
方案二:减少 Master 内存中保存的作业信息:
spark.ui.retainedJobs 500
spark.ui.retainedStages 500
spark.ui.retainedApplications 500
2. Worker 节点假死或失效
Web UI 中发现 Worker 节点消失或处于 dead 状态,任务报错 lost worker。原因与上述类似,Worker 内存中保存了大量 UI 信息,导致 GC 期间与 Master 失去心跳连接。
解决方案
方案一:增加 Worker 内存,在 Worker 节点配置:
export SPARK_DAEMON_MEMORY 2g
方案二:减少 Worker 内存中保存的信息:
spark.worker.ui.retainedExecutors 200
spark.worker.ui.retainedDrivers 200
二、运行时错误处理
1. Shuffle 操作导致的 FetchFailedException
错误特征
错误类型一:缺少输出位置
org.apache.spark.shuffle.MetadataFetchFailedException:
Missing an output location for shuffle 2
错误类型二:Shuffle 拉取失败
org.apache.spark.shuffle.FetchFailedException:
Failed to connect to worker-node-03/192.168.47.215:50268
解决方案
该问题通常出现在大量 Shuffle 操作场景,Task 反复失败重试,最终导致 Application 失败。解决思路如下:
- 增加 Executor 内存配置
- 增加每个 Executor 的 CPU 核心数,保持并行度
推荐配置:
spark.executor.memory 18G
spark.executor.cores 4
spark.cores.max 24
计算公式:
executorCount = spark.cores.max / spark.executor.cores
实际资源配置示例:4核 Executor × 6个 = 24核心,内存总量 108GB。优化后原本运行数小时卡住的任务可在几分钟内完成。
2. Executor 与 Task 丢失
错误特征
类型一:Executor 丢失
WARN TaskSetManager: Lost task 3.0 in stage 1.0 (TID 5, node01.local):
ExecutorLostFailure (executor lost)
类型二:Task 丢失
WARN TaskSetManager: Lost task 70.3 in stage 8.0 (TID 1150, 192.168.47.217):
java.io.IOException: Connection from /192.168.47.217:55483 closed
类型三:各种超时
java.util.concurrent.TimeoutException: Futures timed out after [120 second]
ERROR TransportChannelHandler: Connection to /192.168.47.212:35409
has been quiet for 120000 ms while there are outstanding requests.
Assuming connection is dead; please adjust spark.network.timeout
解决方案
问题由网络抖动或 GC 引起,Worker 或 Executor 未及时接收心跳反馈。调整网络超时配置:
spark.network.timeout 300s
spark.core.connection.ack.wait.timeout 300s
spark.akka.timeout 300s
spark.storage.blockManagerSlaveTimeoutMs 300000
spark.shuffle.io.connectionTimeout 300s
spark.rpc.askTimeout 300s
spark.rpc.lookupTimeout 300s
默认值通常为 120 秒,建议调整为 300 秒(5分钟)或更高。
3. 数据倾斜与任务倾斜
问题描述
- 数据倾斜:部分分区数据量远超其他分区
- 任务倾斜:部分 Task 执行极慢
解决方案
大多数 Task 已完成,个别 Task 长时间无法结束,分为两种情况:
数据倾斜处理:
多数情况由无效数据引起(如 null、空字符串),或异常数据(如某用户登录次数过千万)。无效数据需提前过滤:
sqlContext.sql("SELECT * FROM table WHERE column IS NOT NULL AND column != ''")
原则:多使用 filter 减少实际处理数据量。
任务倾斜处理:
可能由网络 IO、CPU、内存等因素导致。建议启用 Spark 推测执行机制:
spark.speculation true
spark.speculation.interval 100
spark.speculation.quantile 0.75
spark.speculation.multiplier 1.5
4. OOM 内存溢出
错误特征
java.lang.OutOfMemoryError: Java heap space
解决方案
分为 Driver OOM 和 Executor OOM 两种:
Driver OOM: 使用 collect 操作将 Executor 数据汇聚到 Driver 端导致。尽量避免使用 collect。
Executor OOM:
- 增加
spark.executor.memory配置 - 提高任务并行度,将大任务拆分为小任务
5. Task 不可序列化
错误特征
org.apache.spark.SparkException: Job aborted due to stage failure:
Task not serializable: java.io.NotSerializableException: ...
解决方案
在 Worker 中调用 Driver 定义的外部变量时,这些变量未进行序列化传输。问题代码示例:
val config = new Configuration()
dataRDD.map(record => config.process(record)).collect()
三种解决方案:
- 将外部变量移入算子内部,推荐使用 foreachPartition 减少变量创建开销
- 使用
@transient注解标记不需要序列化的变量(如 SparkConf、SparkContext) - 将外部变量封装到可序列化的类中
6. Driver 结果集过大
错误特征
Caused by: org.apache.spark.SparkException:
Job aborted due to stage failure: Total size of serialized
results of 380 tasks (1028.0 MB) is bigger than
spark.driver.maxResultSize (1024.0 MB)
解决方案
每个 Spark Action(如 collect)所有分区的序列化结果总大小受限。避免使用 countByValue、countByKey 等方法,或调大配置:
spark.driver.maxResultSize 2g
7. Task 体积过大
错误特征
WARN TaskSetManager: Stage 198 contains a task of very large size (5953 KB).
The maximum recommended task size is 100 KB.
可能引发错误:
Caused by: java.lang.RuntimeException: Failed to commit task
Caused by: org.apache.spark.executor.CommitDeniedException:
attempt_202106151514_0218_m_000245_0: Not committed
解决方案
Stage 中 Task 过大,通常因 Transform 链过长。解决思路:拆分 Stage,在执行过程中插入 cache 操作切断长链:
val result = data.map(...).filter(...).cache()
result.count() // 触发缓存
result.reduceByKey(...)
8. Driver 未授权提交
Task 完成提交时被 Driver 拒绝,通常与前述 Task 过大问题相关。
9. 环境配置错误
Driver 节点内存不足
Java HotSpot(TM) 64-Bit Server VM warning: INFO:
os::commit_memory(0x0000000680000000, 4294967296, 0) failed;
error='Cannot allocate memory' (errno=12)
解决:将 Driver 部署到内存充足的机器,或调整 spark.driver.memory 参数。
HDFS 空间不足
Caused by: org.apache.hadoop.ipc.RemoteException(java.io.IOException):
File /tmp/spark-history/app-20211228095652-0072.inprogress
could only be replicated to 0 nodes
ERROR LiveListenerBus: Listener EventLoggingListener threw an exception
解决:清理无用数据或增加节点扩展 HDFS 空间。
Spark 与 Hadoop 版本不匹配
java.io.InvalidClassException: org.apache.spark.rdd.RDD;
local class incompatible: stream classdesc serialVersionUID
解决:下载对应 Hadoop 版本的 Spark 或自行编译。
端口占用过多
16/03/16 16:03:17 ERROR SparkUI: Failed to bind SparkUI
java.net.BindException: 地址已在使用: Service 'SparkUI' failed after 16 retries!
解决:提交任务时指定端口
--conf spark.ui.port=54080
中文编码问题
CSV 文件写入 HDFS 出现乱码。JVM 默认使用系统字符集,各节点不一致导致。解决方案:
在 spark-defaults.conf 中配置:
spark.executor.extraJavaOptions -Dfile.encoding=UTF-8
spark.driver.extraJavaOptions -Dfile.encoding=UTF-8
三、Python 相关错误
1. Python 版本过低
java.io.UIException: Cannot run program "python3": error=2, 没有那个文件或目录
解决:确保 Python 版本符合 Spark 要求,升级或创建软链接。
2. Python 权限不足
java.io.IOException: Cannot run program "python3": error=13, 权限不够
解决:检查新节点 Python 安装权限,确保环境变量配置正确。
3. Pickle 序列化失败
TypeError: ('__cinit__() takes exactly 8 positional arguments (11 given)',
<type 'sklearn.tree._tree.Tree'>, ...)
原因:pickle 文件使用新版 scikit-learn 训练,当前环境为旧版本。解决:升级 scikit-learn 并清理旧数据。
4. Python 编码错误
UnicodeEncodeError: 'ascii' codec can't encode characters in position 0-1
解决方案一:
import sys
reload(sys)
sys.setdefaultencoding('utf-8')
解决方案二:
# 报错
str(u'中国')
# 修复
str(u'中国'.encode('utf-8'))
四、性能优化策略
1. Executor 闲置问题
部分 Executor 不执行任务,原因分析:
分区数过少: 每个分区仅由一个 Task 执行,可通过 repartition 调整。不同数据源并发度受限于:
- HDFS:block 数即为分区数
- MySQL:按读取时的分区规则
- Elasticsearch:分片数即分区数
数据本地性副作用:
TaskSetManager 分发任务时计算数据本地性,优先级:
process(同 Executor) → node_local(同节点) → rack_local(同机架) → any(任意节点)
任务执行过快(小于 locality.wait 时间)时,低优先级任务永不启动,造成资源浪费。
判断公式:
currentTime - lastLaunchTime >= localityWait(currentLocalityLevel)
若大量 Executor 闲置,可降低以下参数(可设为 0),默认值均为 3 秒:
spark.locality.wait
spark.locality.wait.process
spark.locality.wait.node
spark.locality.wait.rack
2. Task 连续重试失败
某 Worker 节点故障时,Task 在该 Executor 上持续重试,最终导致 Application 失败。配置失败黑名单:
spark.scheduler.executorTaskBlacklistTime 30000
Task 在某 Executor 失败后会在其他 Executor 启动,该 Executor 进入黑名单 30 秒。
3. 内存优化
Shuffle 数据量大、RDD 缓存少时,可调整以下参数:
spark.storage.memoryFraction:RDD 缓存占比,默认 60%,缓存少时可降低spark.shuffle.memoryFraction:Shuffle 内存占比,默认 20%spark.rdd.compress:是否压缩序列化 RDD,默认 false,可节省空间但消耗 CPU
4. 并行度优化
MySQL 读取并发度:
spark.default.parallelism
Shuffle 时并行度,Standalone 模式下默认为核心数。过大会增加小任务启动开销,过小则大数据量任务执行缓慢。
SQL Shuffle 并行度:
spark.sql.shuffle.partitions
默认 200,过小导致 OOM 和 Executor 丢失,过大则产生大量碎片 Task 增加调度开销。
调整 Map 阶段并行度:
rdd.repartition(partitionNum)
5. Shuffle 优化
SparkSQL Join 优化:采用 Map-side Join 进行关联优化。
6. 磁盘 IO 优化
合理配置本地磁盘目录,使用 SSD 或高速磁盘提升 Shuffle 写性能。
7. 序列化优化
使用 Kryo 序列化替代 Java 默认序列化,提升性能与内存效率。
8. 数据本地性
不同 Cluster Manager 下数据本地性表现存在差异,当 HDFS 数据本地性异常时需排查。
9. 代码层面优化
- 避免使用不必要的 Shuffle 操作
- 合理使用 cache 和 persist
- 减少宽依赖操作
- 使用广播变量优化小表Join