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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 18:06:07