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
相关产品推荐
相关产品推荐

