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

如何在Airflow的on_failure_callback自定义函数中执行Bash脚本

Airflow on_failure_callback 自定义脚本启动实现建议

核心思路

直接在自定义回调函数里用Python的subprocess模块调用本地脚本——因为on_failure_callback是Airflow调度器/worker进程触发的,本身就在Airflow所在服务器上运行,无需额外远程调用,直接执行本地路径即可。

代码实现示例

1. 基础回调函数(执行bash脚本)

import subprocess
from airflow.models import TaskInstance

def run_on_failure_script(context):
    # 获取失败任务的基本信息(可选,用于脚本传参或日志)
    task_instance: TaskInstance = context['task_instance']
    dag_id = task_instance.dag_id
    task_id = task_instance.task_id
    
    # 要执行的脚本绝对路径(务必用绝对路径,避免路径依赖问题)
    script_path = "/opt/airflow/scripts/failure_handler.sh"
    
    try:
        # 执行脚本,可传递任务信息作为参数
        result = subprocess.run(
            [script_path, dag_id, task_id],
            check=True,
            stdout=subprocess.PIPE,
            stderr=subprocess.PIPE,
            text=True
        )
        # 将脚本输出写入Airflow日志
        print(f"Failure script output: {result.stdout}")
    except subprocess.CalledProcessError as e:
        # 捕获脚本执行失败的情况,避免回调自身报错影响Airflow
        print(f"Failed to run failure script: {e.stderr}")

2. 在DAG中使用回调

把这个函数配置到DAG的default_args里,可让所有任务失败时触发;也可以单独给某个任务设置:

from airflow import DAG
from airflow.operators.bash import BashOperator
from datetime import datetime

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2024, 1, 1),
    'on_failure_callback': run_on_failure_script  # 全局生效
}

with DAG('test_dag', default_args=default_args, schedule_interval='@daily') as dag:
    # 示例任务,也可以单独给任务设置on_failure_callback
    bash_task = BashOperator(
        task_id='bash_task',
        bash_command='exit 1',  # 模拟任务失败
        # on_failure_callback=run_on_failure_script  # 单独生效
    )

关键注意事项

  • 权限问题:确保Airflow运行用户(通常是airflow用户)拥有脚本的执行权限,以及脚本读写相关文件的权限。可以用chmod +x /path/to/script.sh给脚本加执行权,必要时调整文件所属用户组。
  • 绝对路径优先:脚本路径、脚本内引用的文件/命令都用绝对路径,避免Airflow运行环境的PATH变量和本地不一致导致找不到文件。
  • 避免阻塞Airflow:如果脚本执行时间较长,改用subprocess.Popen后台执行,防止回调阻塞Airflow的任务调度:
    # 后台执行脚本,不等待完成
    subprocess.Popen([script_path, dag_id, task_id], stdout=subprocess.PIPE, stderr=subprocess.PIPE)
    
  • 日志捕获:通过stdout和stderr捕获脚本输出,写入Airflow日志,方便后续排查脚本执行情况。
  • 环境变量适配:如果脚本依赖特定的环境变量(比如Python虚拟环境、工具路径),可以在回调里先设置环境变量,或者直接在脚本开头指定(比如bash脚本里用source /path/to/venv/bin/activate)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 13:09:56