Azure HDInsight中Python Spark应用无法正常终止问题求助
问题:Azure HDInsight中Spark应用执行完成后陷入无限重试循环
我在集群头节点通过以下shell脚本批量提交Python Spark应用:
spark-submit app1.py > app1.log spark-submit app2.py > app2.log ... spark-submit app<N>.py > app<N>.log
该脚本在本地集群和AWS环境运行正常,但在Azure HDInsight中,Python程序执行完成且预期结果全部生成后,YARN进程始终处于RUNNING状态,必须手动执行yarn application -kill <APP_ID>才能让脚本切换到下一个应用。
日志中反复出现ERROR RawSocketSender错误并循环重试,相关日志片段:
23/02/07 00:42:45 ERROR RawSocketSender [MdsLoggerSenderThread]: org.fluentd.logger.sender.RawSocketSender java.net.SocketException: Broken pipe (Write failed) at java.net.SocketOutputStream.socketWrite0(Native Method) ... at java.lang.Thread.run(Thread.java:750)
问题原因
这个错误源于HDInsight的MDS(监控诊断服务)日志发送器尝试向Fluentd服务发送日志时连接中断,但内置的无限重试逻辑导致日志发送线程无法退出,进而让YARN误认为应用仍在运行——尽管你的Spark任务本身已经完成,数据结果也完全符合预期。
解决办法
临时方案:提交应用时添加配置禁用MDS日志发送
在每个spark-submit命令中加入参数,关闭MDS相关的日志发送功能:
spark-submit --conf spark.mds.logging.enabled=false app1.py > app1.log
如果上述配置无效,也可以尝试关闭DFS和MapReduce的客户端日志发送:
spark-submit --conf spark.hadoop.dfs.client.log.logger=OFF --conf spark.hadoop.mapreduce.client.log.logger=OFF app1.py > app1.log
永久方案:修改集群级Spark配置
如果需要所有应用都避免这个问题,可以直接在HDInsight集群层面修改配置:
- 登录Azure门户,进入你的HDInsight集群管理页面
- 找到「配置」选项,选择Spark服务的配置项
- 添加以下全局配置键值对:
spark.mds.logging.enabled=falsespark.hadoop.dfs.client.log.logger=OFFspark.hadoop.mapreduce.client.log.logger=OFF
- 保存配置并重启Spark服务,后续提交的所有应用都会自动生效
备选方案:在Python应用中显式终止进程
如果不想修改集群配置,也可以在Python应用的末尾添加代码,强制终止进程,避免日志线程无限重试:
import os import signal # 确保所有数据处理完成后,发送终止信号结束进程 os.kill(os.getpid(), signal.SIGTERM)
注意:必须确保任务已经完成所有数据写入操作后再执行这段代码,防止数据丢失。
内容的提问来源于stack exchange,提问作者Rivo RAMAMONJY
相关产品推荐
相关产品推荐

