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

关于Airflow中获取DAG运行状态及成功后触发邮件实现的咨询

关于Airflow中获取DAG运行状态及成功后触发邮件实现的咨询

你好,我来帮你梳理下这段代码里的问题,以及如何正确实现获取DAG整体状态并触发邮件的需求:

一、当前代码存在的问题

  • 上下文对象错误:你在DAG级别设置了on_success_callback=on_success_dag,但在函数里用context.get("task_instance")是不对的——DAG级别的回调函数,上下文(context)里传递的是dag_run对象(代表整个DAG的运行实例),而非单个任务的task_instance。直接取task_instance会导致取不到值甚至抛出异常。
  • 变量名错误:函数里定义了dag = context.get("task_instance").dag_id,但后面SUBJECT里用了dag_id这个未定义的变量,会引发NameError。
  • 未定义的回调函数:DAG的default_args里指定了on_execute_callback、on_failure_callback、on_success_callback为on_trigger,但代码里没有这个函数的定义,会导致运行报错。

二、正确的实现方式

要获取DAG整体的执行状态并在成功时发送邮件,我们应该基于dag_run对象来获取信息,调整后的代码如下:

from airflow import DAG
from airflow.decorators import task
from datetime import datetime
# 假设你有已经配置好的email发送模块
import email

def on_success_dag(context):
    # 获取DAG运行实例
    dag_run = context["dag_run"]
    dag_id = dag_run.dag_id
    dag_state = dag_run.state
    execution_duration = dag_run.duration

    # 构造邮件内容
    SUBJECT = f"Airflow success alert for {dag_id} DAG"
    MSG = f"""
SUCCESSFULLY COMPLETED THE DAG:{dag_id}
DURATION:{execution_duration}
"""

    # 因为是DAG级别的success回调,这里state必然是success,可以省略判断,但保留也没问题
    if dag_state == "success":
        email.send(MSG, SUBJECT)

# DAG定义
with DAG(
    dag_id="testing_email_for_dag",
    schedule_interval=None,
    default_args={
        "start_date": datetime(2023, 6, 20),
        "retries": 0,
        "catchup": False,
        # 如果不需要任务级别的回调,可以移除这几个未定义的参数
        # "on_execute_callback": on_trigger,
        # "on_failure_callback": on_trigger,
        # "on_success_callback": on_trigger,
    },
    on_success_callback=on_success_dag,
    render_template_as_native_obj=True,
    tags=["test"],
) as dag:
    @task
    def print_hello():
        print("hello")

    print_hello()

三、如何获取DAG整体状态

  • 当使用DAG级别的回调函数(on_success_callback、on_failure_callback等)时,直接从context["dag_run"].state就能拿到DAG的整体运行状态,常见的状态值有:success(成功)、failed(失败)、running(运行中)等。
  • 如果是在任务内部或任务级别的回调里想获取整个DAG的状态,可以通过task_instance.dag_run.state来获取,比如在任务函数里:
    @task
    def check_dag_state(**context):
        ti = context["task_instance"]
        dag_state = ti.dag_run.state
        print(f"Current DAG state: {dag_state}")
    

这样调整后,就能正确在DAG成功执行后触发邮件通知,并且准确获取DAG的整体运行状态啦。

备注:内容来源于stack exchange,提问作者Austin Jackson

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.22 10:18:10