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

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事件监听实现

五、验证方法

  1. 手动修改task1.sh,在开头添加exit 1模拟任务失败
  2. 触发DAG运行,观察Airflow UI中task2是否被标记为失败
  3. 登录远程机器,执行ps aux | grep task2.sh确认进程已被终止

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 00:50:04