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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 21:39:00