Airflow DAG的on_failure_callback任务失败时未触发问题排查
我使用Airflow 2.8.2与Python 3.11.8搭建了DAG,配置了on_failure_callback函数send_message_on_dag_fail,但任务失败时该回调未执行且无相关日志。
相关DAG代码如下:
import pandas as pd from datetime import datetime, timedelta from dateutil.relativedelta import relativedelta import pytz import httpx from airflow.decorators import dag, task from airflow.exceptions import AirflowSkipException from airflow.models import Variable def get_current_period(date: datetime.date = None): tz = pytz.timezone('Europe/Moscow') if date: now = date.date() else: now = datetime.now(tz).date() if now.day <= 15: start_date = (now - relativedelta(months=1)).replace(day=16) end_date = now.replace(day=1) - relativedelta(days=1) else: start_date = now.replace(day=1) end_date = now.replace(day=15) return str(start_date), str(end_date) def send_msg(bot_token: str, chat_id: str, message: str, type:str = 'message' or 'code'): if type == 'message': url = f'https://api.telegram.org/bot{bot_token}/sendMessage?chat_id={chat_id}&text={message}' client = httpx.Client(base_url='https://') return client.post(url) elif type == 'code': url = f'https://api.telegram.org/bot{bot_token}/sendMessage' params = { 'chat_id': chat_id, 'text': message, 'parse_mode': 'Markdown' } client = httpx.Client(base_url='https://') return client.post(url, params=params) def get_xcom_from_context(context, task_id: str, dict_key:str = False): if dict_key: xcom = context['ti'].xcom_pull(task_ids=task_id)[dict_key] else: xcom = context['ti'].xcom_pull(task_ids=task_id) return xcom default_args = { 'owner': 'user', 'depends_on_past': False, 'retries': 1, 'retry_delay': timedelta(minutes=1), 'start_date': datetime(2024, 9, 19) } host = Variable.get('host') database_name = Variable.get('database_name') user_name = Variable.get('user_name') password_for_db = Variable.get('password_for_db') server_host_name = Variable.get('server_host_name') bearer_key = Variable.get('bearer_key') user_key = Variable.get('user_key') sales_plans_url = Variable.get('sales_plans_url') specialization_prices_url = Variable.get('specialization_prices_url') bot_token = Variable.get('bot_token') chat_id = Variable.get('chat_id') def send_message_on_dag_fail(bot_token = bot_token, chat_id = chat_id, **kwargs): context = kwargs log = context['ti'].log log.error('DAG FINISHED WITH ERROR __________________') # this error text easier to find task_id = context['ti'].task_id dag_id = context['dag'].dag_id message = f"Task {task_id} from Dag {dag_id} failed." log.error(message) send_msg(bot_token, chat_id, message, 'message') @dag(default_args=default_args, schedule_interval=None, catchup=False, concurrency=4, on_failure_callback=send_message_on_dag_fail) def dag_get_bonus_and_penaltys_for_staff(): @task def check_time(): tz = pytz.timezone('Europe/Moscow') current_time = datetime.now(tz).time() if current_time >= datetime.strptime("00:00", "%H:%M").time() and current_time <= datetime.strptime("01:00", "%H:%M").time(): return False else: return True @task def get_start_end_dates(bot_token = bot_token, chat_id = chat_id,**kwargs): context = kwargs log = context['ti'].log check_time = get_xcom_from_context(context, 'check_time') log.info('Xcom objects pulled from context') if check_time: start_date, end_date = get_current_period() result = { 'start_date': start_date , 'end_date': end_date } task_id = context['ti'].task_id dag_id = context['dag'].dag_id message = f"TASK {task_id} DAG {dag_id}." log.error(message) send_msg(bot_token, chat_id, message, 'message') return result else: raise AirflowSkipException("Time for database cleaning, skip DAG execution.") ... check_time_task = check_time() get_start_end_dates_task = get_start_end_dates() check_time_task >> get_start_end_dates_task >> ... dag_get_bonus_and_penaltys_for_staff = dag_get_bonus_and_penaltys_for_staff()
补充说明:
- 所有
@task均可正常运行 - 已通过为函数添加
**kwargs处理任务实例(其他方式无效) - 已在
get_start_end_dates任务中验证send_message_on_dag_fail核心逻辑可用,为何DAG失败时该回调未生效?
回调作用范围理解错误:你给DAG配置的
on_failure_callback是DAG级别回调,仅当整个DAG实例被标记为失败时才会触发,而非单个任务失败时触发。如果需要单个任务失败就触发回调,要么给每个任务单独配置,要么把回调放到default_args中(所有继承默认参数的任务都会生效)。回调函数参数传递问题:你的
send_message_on_dag_fail用全局变量作为默认参数,Airflow调度时可能因上下文加载顺序问题导致变量未正确初始化。建议直接在函数内从Variable获取值:def send_message_on_dag_fail(**kwargs): bot_token = Variable.get('bot_token') chat_id = Variable.get('chat_id') context = kwargs log = context['ti'].log log.error('DAG FINISHED WITH ERROR __________________') task_id = context['ti'].task_id dag_id = context['dag'].dag_id message = f"Task {task_id} from Dag {dag_id} failed." log.error(message) send_msg(bot_token, chat_id, message, 'message')DAG失败判定条件不符:DAG级别回调仅在DAG整体状态为
failed时触发。如果只是单个任务失败,但DAG因分支逻辑等最终状态为success或skipped,回调不会执行。去Airflow UI的DAG运行页面确认DAG整体状态是否为failed。日志查看位置错误:DAG级别的回调日志不会出现在单个任务日志里,而是在Scheduler日志或DAG的
dag_run日志中。直接搜索你加的标记DAG FINISHED WITH ERROR __________________即可找到相关记录。正确的任务失败回调配置方式:如果需要单个任务失败就发消息,推荐把回调放到
default_args中:default_args = { 'owner': 'user', 'depends_on_past': False, 'retries': 1, 'retry_delay': timedelta(minutes=1), 'start_date': datetime(2024, 9, 19), 'on_failure_callback': send_message_on_dag_fail # 添加到默认参数 } # DAG定义时无需再单独指定on_failure_callback @dag(default_args=default_args, schedule_interval=None, catchup=False, concurrency=4) def dag_get_bonus_and_penaltys_for_staff(): # ... 任务定义
内容的提问来源于stack exchange,提问作者John Doe

