如何实现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
相关产品推荐
相关产品推荐

