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

Airflow跨DAG调用自定义邮件回调时如何传递收件人列表参数

Airflow 自定义失败回调传递差异化收件人方案

现有问题

当前自定义告警函数custom_failure_email内硬编码了统一收件人列表,所有绑定该回调的DAG失败后都会给同一批人发信,无法满足不同DAG、不同任务配置不同告警接收人的需求。

推荐方案:通过context读取DAG/任务自定义参数

Airflow触发失败回调时,传入的context对象会携带当前DAG、任务实例的全量配置,你只需要在DAG/任务的params字段中配置专属收件人列表,回调函数内直接从context读取即可,支持DAG级、任务级多粒度配置,优先级为「任务配置 > DAG配置 > 默认兜底配置」。

第一步:改造自定义告警函数

修改custom_alert.py代码,修复原代码中邮件正文未闭合三引号的语法问题,增加从context读取收件人列表的逻辑:

from airflow.utils.email import send_email

def custom_failure_email(context, **kwargs):
    """发送自定义失败告警邮件"""
    task_instance = context.get('task_instance')
    dag_id = task_instance.dag_id
    # 按优先级读取收件人列表
    alert_recipients = task_instance.params.get(
        'alert_recipients',
        context['dag'].params.get(
            'alert_recipients',
            ['default-alert@company.com'] # 全局兜底收件人,可替换为原有默认列表
        )
    )
    subject = f"[ActionReq]-dag failure-{dag_id}"
    body = """Hi Team,<br><br>
            <b style="font-size:15px;color:red;">Airflow job on error, please find details below.</b>
            Thank you!,<br>
            """
    for email in alert_recipients:
        send_email(to=email, subject=subject, html_content=body)

第二步:在DAG中配置对应收件人

在父DAG文件email.py中,根据需要在不同层级配置收件人即可:

from datetime import datetime
from airflow import DAG
from airflow.operators.python import PythonOperator
from custom_alert import custom_failure_email

# 可在default_args中配置该DAG默认告警收件人
default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'email': ['noreply@gmail.com'],
    'email_on_failure': False,
    'on_failure_callback': custom_failure_email,
    'params': {
        'alert_recipients': ['dag-owner@company.com', 'ops@company.com']
    }
}

with DAG(
    dag_id='business_dag_01',
    default_args=default_args,
    start_date=datetime(2024, 1, 1),
    schedule_interval='@daily',
    # 若DAG下所有任务收件人一致,也可直接在DAG初始化时配置params
    # params={'alert_recipients': ['dag-level-recipient@company.com']}
) as dag:
    def demo_task_func():
        print("task running")

    # 普通任务未单独配置params时,继承default_args/DAG层级的收件人
    normal_task = PythonOperator(
        task_id='normal_task',
        python_callable=demo_task_func
    )

    # 特殊任务可单独配置收件人,覆盖上层配置
    core_task = PythonOperator(
        task_id='core_task',
        python_callable=demo_task_func,
        params={
            'alert_recipients': ['core-task-owner@company.com', 'tech-lead@company.com']
        }
    )

备选方案:使用偏函数预绑定收件人

如果不需要任务级粒度的配置,也可以通过functools.partial给回调函数预绑定当前DAG的收件人列表,不同DAG绑定不同的收件人参数即可。

第一步:改造告警函数增加入参

# custom_alert.py
from airflow.utils.email import send_email

def custom_failure_email(context, recipient_list, **kwargs):
    """发送自定义失败告警邮件"""
    dag_id = context.get('task_instance').dag_id
    subject = f"[ActionReq]-dag failure-{dag_id}"
    body = """Hi Team,<br><br>
            <b style="font-size:15px;color:red;">Airflow job on error, please find details below.</b>
            Thank you!,<br>
            """
    for email in recipient_list:
        send_email(to=email, subject=subject, html_content=body)

第二步:DAG中绑定带参数的回调

from functools import partial
from datetime import datetime
from airflow import DAG
from custom_alert import custom_failure_email

# 为当前DAG绑定专属收件人列表
current_dag_alert = partial(
    custom_failure_email,
    recipient_list=['aditya.dhanraj@gmail.com', 'dag02-recipient@company.com']
)

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'email': ['noreply@gmail.com'],
    'email_on_failure': False,
    'on_failure_callback': current_dag_alert # 传入预绑定参数的回调
}

# 后续DAG初始化逻辑正常编写即可

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 23:42:27