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

Airflow v2.2.5中DAG超时发送邮件通知的实现求助

实现Airflow v2.2.5 DAG超时邮件通知

前提确认

确保你已完成Airflow的邮件基础配置(airflow.cfg中正确设置SMTP相关参数,如smtp_host、smtp_user、smtp_password等),由于你已实现任务级失败通知,这部分配置应该已完成。

步骤1:为DAG设置超时阈值

在DAG定义中添加dagrun_timeout参数,指定DAG的最大允许运行时长,超出该时间后DAG会被标记为失败:

from datetime import datetime, timedelta
from airflow import DAG

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2023, 1, 1),
    # 其他默认参数(如重试设置等)
}

with DAG(
    'your_target_dag_id',
    default_args=default_args,
    description='业务DAG描述',
    schedule_interval='@daily',
    # 示例:设置DAG超时为2小时
    dagrun_timeout=timedelta(hours=2),
    catchup=False
) as dag:
    # 此处编写任务定义(如PythonOperator、BashOperator等)

步骤2:自定义DAG失败回调函数发送邮件

DAG超时后会触发失败状态,通过on_failure_callback绑定自定义函数,实现超时邮件通知:

from airflow.utils.email import send_email

def dag_timeout_alert(context):
    dag_run = context.get('dag_run')
    subject = f"[Airflow] DAG {dag_run.dag_id} 超时失败"
    html_content = f"""
    <p>DAG <b>{dag_run.dag_id}</b> 运行超出设定时长,已被标记为失败。</p>
    <p>DAG运行ID: {dag_run.run_id}</p>
    <p>开始时间: {dag_run.start_date}</p>
    <p>设定超时时长: {dag_run.dag.dagrun_timeout}</p>
    """
    # 替换为实际收件邮箱
    send_email(to=['your_notify_email@xxx.com'], subject=subject, html_content=html_content)

# 在DAG定义中添加回调参数
with DAG(
    'your_target_dag_id',
    default_args=default_args,
    description='业务DAG描述',
    schedule_interval='@daily',
    dagrun_timeout=timedelta(hours=2),
    on_failure_callback=dag_timeout_alert,
    catchup=False
) as dag:
    # 任务定义...

可选:仅针对超时失败发送通知

如果需要区分超时失败与其他原因导致的失败,可在回调函数中增加判断逻辑:

def dag_timeout_alert(context):
    dag_run = context.get('dag_run')
    # 校验是否为超时导致的失败
    if dag_run.state == 'failed' and (dag_run.end_date - dag_run.start_date) > dag_run.dag.dagrun_timeout:
        subject = f"[Airflow] DAG {dag_run.dag_id} 超时失败"
        html_content = f"""
        <p>DAG <b>{dag_run.dag_id}</b> 运行超出设定时长,已被标记为失败。</p>
        <p>DAG运行ID: {dag_run.run_id}</p>
        <p>开始时间: {dag_run.start_date}</p>
        <p>结束时间: {dag_run.end_date}</p>
        <p>设定超时时长: {dag_run.dag.dagrun_timeout}</p>
        """
        send_email(to=['your_notify_email@xxx.com'], subject=subject, html_content=html_content)

验证方式

手动触发目标DAG后,通过延长任务运行时间(如给任务添加time.sleep())让DAG超出dagrun_timeout设定的时长,等待DAG被标记为失败后,检查收件邮箱是否收到通知邮件。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 21:30:48