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

如何实现Airflow标记任务失败时终止Yarn上的Spark应用进程?

解决Airflow任务标记失败时自动终止Yarn应用的方案

方法1:利用任务失败回调(on_failure_callback)

  • 任务执行时,先获取Yarn应用ID并存入Airflow的XCom。比如Spark任务可通过spark.sparkContext.applicationId拿到ID,再用ti.xcom_push(key='yarn_app_id', value=app_id)存储。
  • 定义失败回调函数,从XCom取出应用ID后调用Yarn命令终止:
    def kill_yarn_app(context):
        ti = context['task_instance']
        app_id = ti.xcom_pull(key='yarn_app_id', task_ids='your_task_id')
        if app_id:
            import subprocess
            subprocess.run(['yarn', 'application', '-kill', app_id], check=True)
    
  • 在任务定义里指定on_failure_callback=kill_yarn_app即可。

方法2:自定义Operator重写on_kill方法

  • 继承Airflow现有Operator(比如SparkSubmitOperator),重写on_kill方法加入终止逻辑:
    from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
    import subprocess
    
    class SparkSubmitWithYarnKillOperator(SparkSubmitOperator):
        def on_kill(self):
            # 提前通过日志解析或变量存储获取Yarn应用ID
            app_id = self._get_yarn_app_id()
            if app_id:
                subprocess.run(['yarn', 'application', '-kill', app_id], check=True)
    
  • 使用这个自定义Operator提交任务,任务被标记失败或手动终止时,on_kill会自动触发执行。

方法3:在任务代码中捕获终止信号

  • 在任务Python代码里监听Airflow终止任务时发送的SIGTERM信号,在处理函数里终止Yarn应用:
    import signal
    import subprocess
    
    # 提前保存好Yarn应用ID
    yarn_app_id = ""
    
    def handle_sigterm(signum, frame):
        global yarn_app_id
        if yarn_app_id:
            subprocess.run(['yarn', 'application', '-kill', yarn_app_id], check=True)
        raise SystemExit("任务终止,已同步停止Yarn应用")
    
    # 注册信号处理
    signal.signal(signal.SIGTERM, handle_sigterm)
    
    # 执行任务逻辑,比如启动Spark作业并获取app_id
    # yarn_app_id = spark.sparkContext.applicationId
    

注意事项

  • 确保Airflow Worker节点能正常执行yarn命令,需配置好Yarn环境变量。
  • 若是Spark任务,直接调用spark.sparkContext.stop()也能终止作业,效果和Yarn kill命令一致。

内容的提问来源于stack exchange,提问作者dmitry

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 02:41:26