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.

请求协助解决该问题。
内容的提问来源于stack exchange,提问作者Oyindamola Victor
相关产品推荐
相关产品推荐

