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

如何在Airflow DAG层级失败回调中获取失败任务的task_id

问题根因

你遇到的三个方案失效问题,都是Airflow 2.1.x版本的固有特性和已知bug导致的:

  • 方案1、2拿到随机task_id:DAG层级的on_failure_callback是由第一个将DAG状态标记为失败的任务触发的,并发任务场景下,哪个任务先完成状态上报、先触发DAG状态流转完全不确定,因此从context里直接拿ti或者task_instance_key_str拿到的只是触发回调的那个随机任务,不是真正的失败任务,更不是最早失败的任务。
  • 方案3回调静默失败:Airflow 2.1.1版本中,DAG层级回调传入的dag_run对象未绑定数据库会话,调用get_task_instances()时会触发懒加载会话异常,且该版本调度器没有捕获DAG回调的外层异常打印日志,因此会表现为回调完全不执行、无任何报错输出。

可行解决方案

以下两种方案均适配你当前使用的Airflow 2.1.1版本,且天然保证多任务失败时只发送一条通知,不会造成消息骚扰。

方案1:修复回调逻辑,直接查询元数据库获取失败任务

绕开context中未绑定会话的dag_run对象,直接通过Airflow ORM查询当前DAG运行下所有状态为failed的任务实例,按失败时间排序取最早失败的任务即可,不会出现静默失败问题。

from airflow.models import TaskInstance

def fail_notifier(context):
    current_dag_id = context["dag"].dag_id
    current_run_id = context["run_id"]

    # 直接查询元数据库,精确匹配当前DAG运行的失败任务
    failed_task_instances = TaskInstance.query.filter(
        TaskInstance.dag_id == current_dag_id,
        TaskInstance.run_id == current_run_id,
        TaskInstance.state == "failed"
    ).order_by(TaskInstance.end_date.asc()).all()

    if failed_task_instances:
        first_failed_ti = failed_task_instances[0]
        first_failed_task_id = first_failed_ti.task_id
        all_failed_task_ids = [ti.task_id for ti in failed_task_instances]
        alert_log_url = first_failed_ti.log_url
    else:
        # 兜底逻辑:适配DAG超时等无明确失败任务的场景
        first_failed_task_id = context["ti"].task_id
        all_failed_task_ids = [first_failed_task_id]
        alert_log_url = context["ti"].log_url

    # 此处编写你的Microsoft Teams消息发送逻辑即可
    # 可以将所有失败任务ID、首个失败任务ID、日志链接等信息拼入通知卡片
    send_teams_msg(
        dag_id=current_dag_id,
        first_failed_task=first_failed_task_id,
        all_failed_tasks=all_failed_task_ids,
        log_url=alert_log_url
    )

注意:你配置了max_active_runs=1,加上run_id作为过滤条件可以100%精确匹配当前运行的DAG,不会查到其他运行记录的任务。DAG层级的失败回调只会在DAG状态从running转为failed时触发一次,无论多少个任务并发失败都不会重复执行,完全不会出现多通知骚扰> 骚扰问题。

方案2:将告警逻辑下沉到普通任务(稳定性最高)

如果不想直接操作ORM,可以去掉DAG层级的失败回调,新增一个触发规则为all_done的收尾任务,在任务内判断DAG运行状态、收集失败任务发送通知。该逻辑运行在普通任务进程中,完全不受回调bug影响,报错会直接打印在任务日志中,排查成本极低。

def failure_alert_task(**context):
    dag_run = context["dag_run"]
    # DAG运行成功则直接退出,不发告警
    if dag_run.state != "failed":
        return
    # 任务实例内的dag_run已绑定数据库会话,get_task_instances可正常调用
    failed_tis = dag_run.get_task_instances(state="failed")
    failed_task_ids = [ti.task_id for ti in failed_tis]
    first_failed_ti = sorted(failed_tis, key=lambda x: x.end_date)[0]
    # 此处编写Teams通知发送逻辑
    send_teams_msg(...)

# DAG内新增告警任务
send_alert = PythonOperator(
    task_id="send_failure_alert",
    python_callable=failure_alert_task,
    trigger_rule="all_done",  # 保证前面所有任务成/败/跳过时该任务都能执行
)

# 将告警任务设为所有业务分支的最终下游
chain(final_task, send_alert)

踩坑提醒

Airflow 2.3之前的版本,DAG层级回调传入的多个内置对象都存在懒加载未绑定会话的问题,不要直接调用这类对象的数据库查询方法;涉及状态查询、元数据读取的逻辑,要么直接走ORM查询,要么下沉到普通任务中执行,稳定性远高于写在回调里。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 08:18:27