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

