当前位置:首页 > 技术 > 正文内容

Spark故障排除与性能优化实战指南

访客 技术 2026年7月21日 3

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()

三种解决方案:

  1. 将外部变量移入算子内部,推荐使用 foreachPartition 减少变量创建开销
  2. 使用 @transient 注解标记不需要序列化的变量(如 SparkConf、SparkContext)
  3. 将外部变量封装到可序列化的类中

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
标签: Spark大数据

相关文章

Linux crontab 详解

1) crontab 是什么cron 是 Linux 的定时任务守护进程;crontab 是用来编辑/查看“按时间周期执行命令”的表(cron table)。常见两类:用户 crontab:每个用户一份(crontab -e 编辑)系统级 crontab / cron.d:可指定执行用户(/etc/crontab、/etc/cron.d/*)2) crontab 时间...

富文本里可以允许的 HTML 属性

一、所有标签默认允许的安全属性(极少)class        (可选)id           (通常建议禁用)title️ 注意:id 容易被滥用做锚点注入,很多系统直接禁用class 允许的话最好只允许固定前缀(如 editor-*)二、a 标签允许属性<a href="" t...

Mac 安装 Node.js 指南

方法一:通过官网安装包(最简单,适合初学者)如果你只是想快速安装并开始使用,这是最直接的方法。访问 Node.js 官网。页面会显示两个版本:LTS (Recommended For Most Users):长期支持版,最稳定。建议选这个。Current:最新特性版,包含最新功能但可能不够稳定。下载 .pkg 安装包并运行。按照安装向导点击“下一步”即可完成。方法二:使用 Homebrew 安装(...

Dom\HTML_NO_DEFAULT_NS 的副作用:自动加闭合标签

在使用Dom\HTMLDocument时,Dom\HTML_NO_DEFAULT_NS 将禁止在解析过程中设置元素的命名空间, 此设置是为了与DOMDocument向后兼容而存在的。当使用它时,已知的一个副作用就是:自动加闭合标签例如 </img> 为什么会这样?当你使用:Dom\HTML_NO_DEFAULT_NS文档会变成 无命名空间模式,此时内部更接近 XML...

Laravel 事件和监听器创建

在 Laravel 中,使用 Artisan 命令创建 Events(事件) 和 Listeners(监听器) 是非常高效的。你可以通过以下几种方式来实现:1. 手动创建单个 Event如果你只想创建一个事件类,可以使用 make:event 命令:Bashphp artisan make:event UserRegistered执行后,文件将生成在 app/Even...

自定义域名解析神器 dnsmasq

什么是 dnsmasq?dnsmasq 是一个轻量级、功能强大的网络服务工具,专为小型和中等规模网络设计。它是一个综合的网络基础设施解决方案[1]。dnsmasq 能做什么?功能说明应用场景DNS 转发与缓存将 DNS 查询转发到上游服务器(ISP、Google DNS 等),并在本地缓存结果加快 DNS 查询速度,减少外部 DNS 流量本地 DNS解析本地网络设备的主机名,无需编辑&n...

发表评论

访客

◎欢迎参与讨论,请在这里发表您的看法和观点。