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

Spark 3.3.3 on Yarn集群心跳超时与Executor丢失问题求助

Spark 3.3.3 Yarn集群流处理任务Executor心跳超时/容器终止问题

背景

我查过StackOverflow相关帖子,但都不匹配我的情况——那些帖子里的Spark版本都是5年以前的,我用的是Spark 3.3.3。

我在Yarn为主节点的Apache Spark集群上运行任务,用Jupyter Labs做IDE,启动集群的代码如下:

from pyspark.sql import SparkSession

import getpass
username = getpass.getuser()

spark = SparkSession. \
    builder. \
    config('spark.ui.port', '0'). \
    config('spark.sql.warehouse.dir', f'/user/{username}/warehouse'). \
    enableHiveSupport(). \
    appName(f'{username} | Python - Kafka and Spark Integration'). \
    master('yarn'). \
    getOrCreate()

任务描述

我的任务是从HDFS读取24个约80MB的JSON文件流,按created_year、created_month、created_dayofmonth三列分区后,以parquet格式写入HDFS另一目录,执行代码:

file_df. \
    writeStream. \
    partitionBy('created_year', 'created_month', 'created_dayofmonth'). \
    format('parquet'). \
    option("checkpointLocation", f"/user/{username}/file_df/streaming/checkpoint/file_df"). \
    option("path", f"/user/{username}/file_df/streaming/data/files_parq"). \
    trigger(once=True). \
    start()

错误日志

执行后出现Executor心跳超时、连接重置、容器被终止等错误,任务阶段进度推进后最终中断,具体日志:

23/09/13 21:34:55 WARN ResolveWriteToStream: spark.sql.adaptive.enabled is not supported in streaming DataFrames/Datasets and will be disabled.

<pyspark.sql.streaming.StreamingQuery at 0x7fd22851adc0>

[Stage 3:===============>                                          (5 + 2) / 19]

23/09/13 21:42:56 WARN HeartbeatReceiver: Removing executor 2 with no recent heartbeats: 120168 ms exceeds timeout 120000 ms
23/09/13 21:42:56 ERROR YarnScheduler: Lost executor 2 on sparkde.camp.300123.internal: Executor heartbeat timed out after 120168 ms
23/09/13 21:42:56 WARN TaskSetManager: Lost task 1.0 in stage 3.0 (TID 41) (sparkde.camp.300123.internal executor 2): ExecutorLostFailure (executor 2 exited caused by one of the running tasks) Reason: Executor heartbeat timed out after 120168 ms

[Stage 3:===============>                                          (5 + 1) / 19]

23/09/13 21:42:58 WARN TransportChannelHandler: Exception in connection from /10.172.0.3:41888
java.io.IOException: Connection reset by peer
    at sun.nio.ch.FileDispatcherImpl.read0(Native Method)
    at sun.nio.ch.SocketDispatcher.read(SocketDispatcher.java:39)
    at sun.nio.ch.IOUtil.readIntoNativeBuffer(IOUtil.java:223)
    at sun.nio.ch.IOUtil.read(IOUtil.java:192)
    at sun.nio.ch.SocketChannelImpl.read(SocketChannelImpl.java:379)
    at io.netty.buffer.PooledByteBuf.setBytes(PooledByteBuf.java:258)
    at io.netty.buffer.AbstractByteBuf.writeBytes(AbstractByteBuf.java:1132)
    at io.netty.channel.socket.nio.NioSocketChannel.doReadBytes(NioSocketChannel.java:350)
    at io.netty.channel.nio.AbstractNioByteChannel$NioByteUnsafe.read(AbstractNioByteChannel.java:151)
    at io.netty.channel.nio.NioEventLoop.processSelectedKey(NioEventLoop.java:722)
    at io.netty.channel.nio.NioEventLoop.processSelectedKeysOptimized(NioEventLoop.java:658)
    at io.netty.channel.nio.NioEventLoop.processSelectedKeys(NioEventLoop.java:584)
    at io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:496)
    at io.netty.util.concurrent.SingleThreadEventExecutor$4.run(SingleThreadEventExecutor.java:986)
    at io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74)
    at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30)
    at java.lang.Thread.run(Thread.java:750)

[Stage 3:========================>                                 (8 + 2) / 19]

23/09/13 21:45:55 ERROR YarnScheduler: Lost executor 3 on sparkde.camp.300123.internal: Container from a bad node: container_1694638458719_0001_01_000004 on host: sparkde.camp.300123.internal. Exit status: 143. Diagnostics: [2023-09-13 21:45:55.150]Container killed on request. Exit code is 143
[2023-09-13 21:45:55.151]Container exited with a non-zero exit code 143. 
[2023-09-13 21:45:55.151]Killed by external signal
.
23/09/13 21:45:55 WARN TaskSetManager: Lost task 1.1 in stage 3.0 (TID 47) (sparkde.camp.300123.internal executor 3): ExecutorLostFailure (executor 3 exited caused by one of the running tasks) Reason: Container from a bad node: container_1694638458719_0001_01_000004 on host: sparkde.camp.300123.internal. Exit status: 143. Diagnostics: [2023-09-13 21:45:55.150]Container killed on request. Exit code is 143
[2023-09-13 21:45:55.151]Container exited with a non-zero exit code 143. 
[2023-09-13 21:45:55.151]Killed by external signal
.
23/09/13 21:45:55 WARN YarnSchedulerBackend$YarnSchedulerEndpoint: Requesting driver to remove executor 3 for reason Container from a bad node: container_1694638458719_0001_01_000004 on host: sparkde.camp.300123.internal. Exit status: 143. Diagnostics: [2023-09-13 21:45:55.150]Container killed on request. Exit code is 143
[2023-09-13 21:45:55.151]Container exited with a non-zero exit code 143. 
[2023-09-13 21:45:55.151]Killed by external signal

阶段进度持续推进后最终停止,出现以下错误:

[Stage 3:=======================================>                 (13 + 1) / 19]

23/09/13 21:54:46 WARN TaskSetManager: Lost task 14.0 in stage 3.0 (TID 58) (sparkde.camp.300123.internal executor 6): TaskKilled (Stage cancelled)
23/09/13 22:03:05 WARN SparkContext: Executor 1 might already have stopped and can not request thread dump from it.

Spark任务详情截图

请求协助解决该问题。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 22:24:52