Airflow中SSHOperator任务失败时如何终止并发任务?
解决Airflow并发SSHOperator任务的联动终止问题
一、核心问题分析
你之前的方案仅实现了任务失败后的下游触发逻辑,但无法主动终止正在运行的SSHOperator任务;同时SSHOperator默认机制下,任务终止时无法可靠清理远程机器上的进程,导致残留。要解决这个问题,需从「Airflow任务联动终止」和「远程进程优雅退出」两个维度处理。
二、步骤1:实现任务失败时主动终止其他并发任务
利用Airflow的任务失败回调机制,在任一SSHOperator任务失败时,调用Airflow内部API终止另一并发任务。注意需确保Airflow API可被任务访问,且任务拥有足够权限操作其他任务状态。
代码实现
from airflow import DAG from airflow.providers.ssh.operators.ssh import SSHOperator from airflow.exceptions import AirflowException from datetime import datetime import requests # 替换为你的Airflow API地址和管理员Token AIRFLOW_API_URL = "http://localhost:8080/api/v1" AIRFLOW_API_TOKEN = "your_admin_token_here" def terminate_sibling_task(context): """任务失败时终止另一并发任务""" dag_id = context["dag"].dag_id failed_task_id = context["task_instance"].task_id target_task_id = "task2" if failed_task_id == "task1" else "task1" # 调用Airflow API终止目标任务 api_url = f"{AIRFLOW_API_URL}/dags/{dag_id}/dagRuns/{context['dag_run'].run_id}/taskInstances/{target_task_id}" headers = {"Authorization": f"Bearer {AIRFLOW_API_TOKEN}"} payload = {"state": "failed", "note": "Terminated due to sibling task failure"} try: response = requests.patch(api_url, json=payload, headers=headers) response.raise_for_status() print(f"Successfully terminated task {target_task_id}") except Exception as e: print(f"Failed to terminate task {target_task_id}: {str(e)}") raise AirflowException(f"Failed to terminate sibling task: {str(e)}") with DAG( dag_id="ssh_parallel_task", schedule_interval=None, start_date=datetime(2024, 1, 1), catchup=False ) as dag: task1 = SSHOperator( task_id='task1', command="/home/user/task1.sh", cmd_timeout=None, get_pty=True, on_failure_callback=terminate_sibling_task ) task2 = SSHOperator( task_id='task2', command="/home/user/task2.sh", cmd_timeout=None, get_pty=True, on_failure_callback=terminate_sibling_task ) # 两个任务并行执行,无依赖关系 task1 task2
三、步骤2:确保远程进程随SSHOperator任务终止而退出
默认情况下,若远程脚本忽略SIGHUP信号(如使用nohup),即使Airflow任务终止,远程进程仍会残留。需修改远程shell脚本,添加信号处理逻辑,确保连接断开时主动终止所有子进程。
修改远程shell脚本(以task1.sh为例)
#!/bin/bash # 定义信号处理函数,收到终止信号时清理子进程 cleanup_processes() { echo "Received termination signal, cleaning up..." kill -TERM $CHILD_PID exit 1 } # 注册信号捕获:SIGHUP对应SSH连接断开,SIGTERM对应主动终止 trap cleanup_processes SIGHUP SIGTERM SIGINT # 执行业务逻辑,这里模拟长时间运行任务 sleep 300 & CHILD_PID=$! # 等待子进程完成,避免脚本提前退出 wait $CHILD_PID
关键说明
- 开启
get_pty=True时,Airflow终止任务会断开SSH连接,远程进程会收到SIGHUP信号 - 通过
trap命令捕获信号,主动终止脚本启动的子进程 - 避免在脚本中使用
nohup后台运行后直接退出,防止子进程变为孤儿进程
四、替代方案:K8s环境下使用PodOperator
若Airflow部署在Kubernetes集群,推荐用PodOperator替代SSHOperator:
- 任务终止时会直接销毁Pod,远程进程会被强制清理,无残留
- 可通过K8s Pod生命周期钩子实现优雅退出
- 任务间的终止联动可通过Airflow回调或K8s事件监听实现
五、验证方法
- 手动修改task1.sh,在开头添加
exit 1模拟任务失败 - 触发DAG运行,观察Airflow UI中task2是否被标记为失败
- 登录远程机器,执行
ps aux | grep task2.sh确认进程已被终止
内容的提问来源于stack exchange,提问作者sblue
相关产品推荐
相关产品推荐

