Yarn-Client模式Spark任务完成后仍挂起问题排查(PySpark3.4.1)
在Kerberos认证的CDP Hadoop集群中,以Docker容器作为Driver节点,提交Yarn-Client模式的PySpark3.4.1任务时,出现应用生命周期不一致的问题,具体表现如下:
1. 任务提交配置
运行在${CONTAINER_HOST}上的Docker容器作为Driver,因防火墙限制使用自定义端口并完成主机端口转发,提交命令如下:
spark-submit --master yarn --deploy-mode client \ --conf spark.driver.host=${CONTAINER_HOST} \ --conf spark.driver.bindAddress=${IP_OF_DOCKER_CONTAINER} \ --conf spark.ui.port=4000 \ --conf spark.driver.port=5000 \ --conf spark.blockManager.port=6000 \ --name sample_yarnclient_job \ --conf spark.executor.memoryOverhead=4096 \ --conf spark.sql.broadcastTimeout=3600 \ --conf spark.sql.autoBroadcastJoinThreshold=-1 \ --conf spark.executor.metrics.pollingInterval=10000 \ --conf spark.executor.heartbeatInterval=3s \ --conf spark.executor.extraJavaOptions=-XX:+UseG1GC \ --conf spark.network.timeout=50000s \ --conf spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version=2 \ --conf spark.yarn.am.memory=12g \ --conf spark.yarn.am.cores=2 \ --conf spark.yarn.tags=sample_yarnclient_job \ --conf spark.network.timeout=50000s \ --conf spark.yarn.historyServer.allowTracking=true \ --conf spark.ui.filters=org.apache.spark.deploy.yarn.YarnProxyRedirectFilter \ --driver-cores 1 \ --driver-memory 6g \ --executor-memory 6g \ --executor-cores 5 \ --num-executors 24 \ main.py
2. 异常现象
- 长时任务可通过
${CONTAINER_HOST}:4000正常访问Spark UI; df.write完成后Driver进程随机挂起;- Driver日志显示SparkContext已正常停止(exitCode 0),但Application Master(AM)无活动约15分钟后被标记为FAILED;
- YARN日志报错无法连接
${CONTAINER_HOST}:5000:
Uncaught exception: org.apache.spark.SparkException: Exception thrown in awaitResult: at org.apache.spark.util.ThreadUtils$.awaitResult(ThreadUtils.scala:322) ... Caused by: java.net.ConnectException: Connection refused ...
3. Driver代码实现
已通过with语句管理SparkSession,确保资源正常释放:
def main() -> None: spark_conf = {'spark.logConf': 'true'} with create_session('sample_yarnclient_job', **spark_conf): SparkContext.getOrCreate().setLogLevel("ERROR") start_time: datetime = datetime.now() some_df_operation(df) end_time: datetime = datetime.now() LOGGER.info(f'duration : {end_time - start_time}')
已调整多种网络超时、心跳参数但未彻底解决,求问任务完成后AM仍与Driver断开连接并失败的原因及解决方法?
核心原因:Yarn-Client模式下AM与Driver的生命周期协同问题
在Yarn-Client模式中,Driver运行在本地容器,AM仅负责向YARN申请资源并管理Executor,任务完成后AM需要向Driver确认最终状态才能正常退出。你的场景中存在以下几个触发点:
容器端口提前释放/网络链路中断
当SparkContext通过with语句关闭后,Driver进程会快速释放spark.driver.port=5000,而此时AM可能仍在尝试与Driver做最后的状态同步。由于容器环境中端口资源回收较快,或者主机端口转发规则在Driver进程退出后被自动清理,导致AM无法连接目标端口,最终触发连接拒绝错误并被YARN标记为FAILED。SparkContext关闭与AM状态同步的时序差
PySpark中with语句关闭SparkSession时,会立即终止SparkContext并释放所有端口,但YARN AM的状态同步逻辑存在延迟:AM需要等待所有Executor退出、清理临时资源后,才会向Driver发送最终的任务完成确认请求,此时Driver端口已被释放,导致连接失败。Kerberos认证会话提前失效
若Driver的Kerberos Ticket在任务完成后快速过期,AM尝试与Driver建立认证连接时会失败,部分场景下Kerberos认证失败会被包装为连接异常。
针对性解决方案
方案1:延迟Driver端口释放,给AM足够同步时间
在with语句块结束后添加延迟,确保AM完成状态同步后再让Driver进程退出:
def main() -> None: spark_conf = {'spark.logConf': 'true'} with create_session('sample_yarnclient_job', **spark_conf): SparkContext.getOrCreate().setLogLevel("ERROR") start_time: datetime = datetime.now() some_df_operation(df) # 延迟30秒,等待AM完成状态同步 import time time.sleep(30) end_time: datetime = datetime.now() LOGGER.info(f'duration : {end_time - start_time}')
方案2:配置AM等待Driver确认的超时时间
添加Spark配置,延长AM等待Driver响应的超时窗口,确保AM在Driver退出前完成状态同步:
# 在spark-submit中追加以下配置 --conf spark.yarn.am.waitTimeBeforeKill=60000 \ # AM等待Driver确认的超时时间,单位毫秒,设为60秒 --conf spark.yarn.driverCommunicationTimeout=300000 # Driver与AM的通信超时时间,单位毫秒
方案3:确保容器端口转发的持久性
检查Docker端口转发规则,确保5000端口的转发不会随Driver进程退出而立即失效。可通过Docker--restart=always配合后台运行容器,或使用宿主机iptables规则手动配置端口转发,而非依赖容器启动时的自动转发。
方案4:延长Kerberos Ticket有效期
确保Driver容器中的Kerberos Ticket有效期覆盖整个任务周期(包括AM同步的时间窗口):
# 在容器启动时执行,设置Ticket有效期为10小时 kinit -kt /path/to/keytab principal@REALM -l 10h
方案5:强制AM在任务完成后主动退出
添加配置让AM在所有Executor退出后直接标记任务为SUCCEEDED,无需等待Driver确认(仅适用于任务结果已持久化、无需Driver额外处理的场景):
--conf spark.yarn.am.exitOnFinish=true
内容的提问来源于stack exchange,提问作者StrangerThinks

