Spark/Scala场景下处理YARN终止事件实现退出前资源清理与状态保存
可直接捕获YARN终止信号的实现方法
YARN终止应用的逻辑和POSIX信号规范完全兼容,完全可以通过捕获信号实现退出前的收尾操作,具体适配方式分部署模式区别处理:
- Cluster部署模式:Spark驱动运行在YARN的ApplicationMaster(AM)容器内,执行
yarn kill <应用ID>命令时,YARN NodeManager会先向AM容器发送SIGTERM信号,等待默认10秒(可通过YARN配置项yarn.nm.container-executor.sleep-delay-before-sigkill调整时长)后再发送SIGKILL强制终止进程。你只需要在驱动代码中注册JVM关闭钩子/语言层面的退出回调即可在SIGTERM触发时执行收尾逻辑:// Scala 驱动示例:注册关闭钩子 sys.addShutdownHook { // 自定义收尾逻辑:保存作业ID、持久化运行状态、清理临时HDFS文件等 println("触发退出前收尾操作") }# PySpark 驱动示例:注册退出回调 import atexit @atexit.register def clean_up(): # 自定义收尾逻辑 print("触发退出前收尾操作") - Client部署模式:Spark驱动运行在提交spark-submit的本地节点,
yarn kill只会终止YARN侧的AM和执行器容器,不会给本地驱动进程发送信号,需要你在执行yarn kill的同时给本地驱动进程发送SIGTERM信号,才能触发本地的收尾逻辑。
注意:禁止使用
yarn kill -signal SIGKILL <应用ID>终止应用,该信号会被内核直接强制终止进程,无法被捕获,不会触发任何收尾逻辑。
替代yarn kill的优雅终止方案
如果不方便修改驱动代码,或者需要更细粒度的终止控制,可以使用以下方案:
- 自定义终止接口:在驱动中启动一个轻量HTTP服务,收到终止请求时优先执行所有收尾逻辑,再主动调用
sc.stop()停止Spark上下文并退出进程,YARN会自动识别应用完成状态并回收所有关联资源。 - 单作业粒度终止:如果只需要停止应用下的某一个作业而不是终止整个Spark应用,可以调用Spark原生REST接口
POST /v1/applications/<应用ID>/jobs/<作业ID>/kill终止指定作业,不会影响其他运行中的作业,也不需要停止整个应用。 - 批量优雅终止脚本:提交应用时给Spark应用打上自定义标签,需要批量终止时先通过
yarn application -list -appTags <自定义标签>获取所有关联应用ID,遍历给每个应用发送SIGTERM,等待预设的收尾时长后再检查应用状态,对未退出的应用再发送SIGKILL强制终止。
注意事项
- 如果配置了
spark.yarn.maxAppAttempts > 1,YARN会在应用异常退出时自动重试,需要在收尾逻辑中主动标记应用为最终失败状态,避免不必要的重试。 - 关闭钩子的执行时长不要超过YARN配置的
SIGTERM到SIGKILL的等待时长,否则逻辑未执行完成进程就会被强制终止,建议提前调整该配置到大于收尾逻辑的最大执行时间。
内容的提问来源于stack exchange,提问作者tomoyo255
相关产品推荐
相关产品推荐

