You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark Worker未手动终止却收到SIGTERM且报Java堆内存溢出求助

Spark Worker意外收到SIGTERM信号(伴随OOM异常)的排查与解决方案

问题场景

  • 集群架构:无外部集群管理器,采用主备Master节点+多Worker节点的cluster部署模式
  • Spark版本:2.2.1
  • 异常现象:Worker进程未执行手动终止操作,却收到SIGTERM信号并启动shutdown流程,日志中先出现**Java堆内存溢出(OutOfMemoryError)**异常

异常日志片段

Exception in thread "dispatcher-event-loop-9" java.lang.OutOfMemoryError: Java heap space
at java.util.zip.ZipCoder.getBytes(ZipCoder.java:80)
at java.util.zip.ZipFile.getEntry(ZipFile.java:310)
at java.util.jar.JarFile.getEntry(JarFile.java:240)
at java.util.jar.JarFile.getJarEntry(JarFile.java:223)
at sun.misc.URLClassPath$JarLoader.getResource(URLClassPath.java:1005)
at sun.misc.URLClassPath.getResource(URLClassPath.java:212)
at java.net.URLClassLoader$1.run(URLClassLoader.java:365)
at java.net.URLClassLoader$1.run(URLClassLoader.java:362)
at java.security.AccessController.doPrivileged(Native Method)
at java.net.URLClassLoader.findClass(URLClassLoader.java:361)
at java.lang.ClassLoader.loadClass(ClassLoader.java:424)
at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:331)
at java.lang.ClassLoader.loadClass(ClassLoader.java:357)
at org.apache.spark.rpc.netty.Inbox.safelyCall(Inbox.scala:206)
at org.apache.spark.rpc.netty.Inbox.process(Inbox.scala:101)
at org.apache.spark.rpc.netty.Dispatcher$MessageLoop.run(Dispatcher.scala:216)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
at java.lang.Thread.run(Thread.java:745)
java.lang.OutOfMemoryError: Java heap space
at java.lang.Class.getDeclaredMethods0(Native Method)
at java.lang.Class.privateGetDeclaredMethods(Class.java:2701)
at java.lang.Class.getDeclaredMethod(Class.java:2128)
at java.io.ObjectStreamClass.getPrivateMethod(ObjectStreamClass.java:1431)
at java.io.ObjectStreamClass.access$1700(ObjectStreamClass.java:72)
at java.io.ObjectStreamClass$2.run(ObjectStreamClass.java:494)
at java.io.ObjectStreamClass$2.run(ObjectStreamClass.java:468)
at java.security.AccessController.doPrivileged(Native Method)
at java.io.ObjectStreamClass.<init>(ObjectStreamClass.java:468)
at java.io.ObjectStreamClass.lookup(ObjectStreamClass.java:365)
at java.io.ObjectStreamClass.initNonProxy(ObjectStreamClass.java:602)
at java.io.ObjectInputStream.readNonProxyDesc(ObjectInputStream.java:1623)
at java.io.ObjectInputStream.readClassDesc(ObjectInputStream.java:1518)
at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:1774)
at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1351)
at java.io.ObjectInputStream.readObject(ObjectInputStream.java:371)
at org.apache.spark.serializer.JavaDeserializationStream.readObject(JavaSerializer.scala:75)
at org.apache.spark.serializer.JavaSerializerInstance.deserialize(JavaSerializer.scala:108)
at org.apache.spark.rpc.netty.NettyRpcEnv$$anonfun$deserialize$1$$anonfun$apply$1.apply(NettyRpcEnv.scala:270)
at scala.util.DynamicVariable.withValue(DynamicVariable.scala:58)
at org.apache.spark.rpc.netty.NettyRpcEnv.deserialize(NettyRpcEnv.scala:319)
at org.apache.spark.rpc.netty.NettyRpcEnv$$anonfun$deserialize$1.apply(NettyRpcEnv.scala:269)
at scala.util.DynamicVariable.withValue(DynamicVariable.scala:58)
at org.apache.spark.rpc.netty.NettyRpcEnv.deserialize(NettyRpcEnv.scala:268)
at org.apache.spark.rpc.netty.RequestMessage$.apply(NettyRpcEnv.scala:603)
at org.apache.spark.rpc.netty.NettyRpcHandler.internalReceive(NettyRpcEnv.scala:654)
at org.apache.spark.rpc.netty.NettyRpcHandler.receive(NettyRpcEnv.scala:646)
at org.apache.spark.network.server.TransportRequestHandler.processOneWayMessage(TransportRequestHandler.java:178)
at org.apache.spark.network.server.TransportRequestHandler.handle(TransportRequestHandler.java:107)
at org.apache.spark.network.server.TransportChannelHandler.channelRead(TransportChannelHandler.java:118)
at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:357)
at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:343)
22/12/24 00:46:43 ERROR Worker: RECEIVED SIGNAL TERM
22/12/24 00:46:51 INFO ExecutorRunner: Killing process!
22/12/24 00:46:51 INFO ExecutorRunner: Killing process!
22/12/24 00:46:51 INFO ExecutorRunner: Killing process!
22/12/24 00:46:51 INFO ExecutorRunner: Killing process!
.....
22/12/24 00:46:51 INFO ShutdownHookManager: Shutdown hook called

原因排查

  1. 系统OOM Killer触发终止:Linux系统中,当节点内存资源耗尽时,内核的OOM Killer会选择占用内存较高的进程发送SIGTERM信号终止,以释放系统内存。日志中Worker先出现堆内存溢出,随后收到SIGTERM,这是最直接的触发原因。
  2. RPC通信线程崩溃引发连锁反应:Worker的RPC处理线程(dispatcher-event-loop)因OOM崩溃,导致无法响应Master的心跳或指令,Master可能标记Worker为失效,间接触发Worker的shutdown流程(Spark 2.2.1中Master不会主动发送SIGTERM,但Worker内部线程崩溃会触发自我终止逻辑)。
  3. Worker内存配置不足:默认的Worker堆内存配置无法支撑当前作业的资源需求,导致内存持续占用直至溢出。

解决方案

1. 调整Worker JVM内存参数

修改Spark配置文件spark-env.sh,增加或调整以下参数:

# 设置Worker堆内存大小,根据节点硬件配置调整,例如节点内存16G时设置8G
SPARK_WORKER_OPTS="-Xmx8g -XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=/var/log/spark/worker_heapdump.hprof"

添加HeapDump参数可在OOM时生成堆转储文件,便于后续分析内存泄漏点。

2. 优化Spark作业内存占用

  • 切换序列化方式:将默认的Java序列化改为Kryo序列化,减少序列化后的内存占用,在spark-defaults.conf中配置:
    spark.serializer org.apache.spark.serializer.KryoSerializer
    spark.kryo.registrationRequired true
    
  • 避免数据倾斜:检查作业中是否存在数据倾斜,通过加盐、分区调整等方式分散数据压力。
  • 合理设置缓存策略:避免无限制缓存RDD,使用persist(StorageLevel.MEMORY_AND_DISK_SER)替代默认的内存缓存,减少内存占用。

3. 系统层面OOM防护

  • 调整OOM Killer优先级:为Spark Worker进程设置较低的OOM分数,降低被选中终止的概率,执行命令:
    echo -1000 > /proc/$(pgrep -f spark.Worker)/oom_score_adj
    
  • 配置内存过载策略:修改/etc/sysctl.conf调整系统内存参数:
    vm.overcommit_memory=1
    vm.panic_on_oom=0
    
    执行sysctl -p使配置生效,此为应急方案,需配合内存根本优化使用。

4. 升级Spark版本

Spark 2.2.1为老旧版本,后续3.x版本对内存管理、RPC通信机制有大幅优化,能有效降低OOM及异常终止的概率。

5. 监控内存状态

部署监控工具实时监控Worker节点的堆内存、系统内存使用情况,提前预警内存不足问题,及时调整配置或优化作业。

内容的提问来源于stack exchange,提问作者Calypso

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.29 00:55:05