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
原因排查
- 系统OOM Killer触发终止:Linux系统中,当节点内存资源耗尽时,内核的OOM Killer会选择占用内存较高的进程发送SIGTERM信号终止,以释放系统内存。日志中Worker先出现堆内存溢出,随后收到SIGTERM,这是最直接的触发原因。
- RPC通信线程崩溃引发连锁反应:Worker的RPC处理线程(dispatcher-event-loop)因OOM崩溃,导致无法响应Master的心跳或指令,Master可能标记Worker为失效,间接触发Worker的shutdown流程(Spark 2.2.1中Master不会主动发送SIGTERM,但Worker内部线程崩溃会触发自我终止逻辑)。
- 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=0sysctl -p使配置生效,此为应急方案,需配合内存根本优化使用。
4. 升级Spark版本
Spark 2.2.1为老旧版本,后续3.x版本对内存管理、RPC通信机制有大幅优化,能有效降低OOM及异常终止的概率。
5. 监控内存状态
部署监控工具实时监控Worker节点的堆内存、系统内存使用情况,提前预警内存不足问题,及时调整配置或优化作业。
内容的提问来源于stack exchange,提问作者Calypso
相关产品推荐
相关产品推荐

