Airflow中使用Jinja模板引用UI变量配置任务失败通知邮件收件人失败问题
解决Airflow任务失败通知邮件中Jinja模板未生效的问题
看起来你踩了Airflow里一个容易忽略的坑:default_args中的email参数并不支持Jinja模板渲染,所以你写的{{ var.json.my_dag.email_on_failure_list}}被直接当成了一个无效的邮箱字符串,才会出现"Invalid address"的错误。
问题根源
Airflow的Jinja模板渲染主要针对任务执行阶段的动态参数(比如bash_command、op_kwargs这类和任务运行相关的参数),而default_args是在DAG解析阶段就被处理的,这个阶段Airflow不会解析其中的Jinja模板语法,你的模板字符串就原封不动地传给了邮件发送模块,自然会被判定为无效地址。
解决方案
这里有两种实用的解决方式,你可以根据场景选择:
方法一:自定义失败回调函数(推荐)
这种方式既避开了在顶层代码直接调用Variable.get()的性能问题,又能灵活控制邮件内容,是Airflow官方更推荐的做法。
from airflow.utils.email import send_email from airflow.models import Variable from datetime import datetime from airflow import DAG from airflow.operators.dummy import DummyOperator def failure_notify_callback(context): # 在任务执行阶段动态获取变量,不会影响DAG解析效率 dag_config = Variable.get("my_dag", deserialize_json=True) email_recipients = dag_config["email_on_failure_list"] # 自定义邮件主题和内容,还能加入任务日志链接等信息 task_instance = context["task_instance"] subject = f"⚠️ 任务失败提醒:{task_instance.task_id}(DAG: {context['dag'].dag_id})" html_content = f""" <h3>任务失败告警</h3> <p>任务名称:{task_instance.task_id}</p> <p>DAG名称:{context['dag'].dag_id}</p> <p>执行时间:{context['execution_date']}</p> <p>日志链接:<a href="{task_instance.log_url}">点击查看详情</a></p> """ # 发送邮件 send_email(to=email_recipients, subject=subject, html_content=html_content) default_args = { 'depends_on_past': False, 'start_date': datetime(2020, 9, 1), 'on_failure_callback': failure_notify_callback # 替换原有的email相关参数 } with DAG('my_dag', default_args=default_args, schedule_interval='@daily') as dag: dummy_task = DummyOperator(task_id='dummy_task')
方法二:DAG解析阶段动态获取变量
如果你更习惯用email_on_failure参数,可以在DAG解析时直接获取变量值(注意:这种方式会在DAG解析时请求数据库,若你有大量DAG可能影响性能,建议加上异常处理):
from airflow.models import Variable from datetime import datetime from airflow import DAG from airflow.operators.dummy import DummyOperator # 处理变量不存在或键缺失的情况,避免DAG解析失败 try: dag_config = Variable.get("my_dag", deserialize_json=True) email_recipients = dag_config["email_on_failure_list"] except (Variable.DoesNotExist, KeyError): email_recipients = ["admin@example.com"] # 设置默认收件人 default_args = { 'depends_on_past': False, 'start_date': datetime(2020, 9, 1), 'email_on_failure': True, 'email': email_recipients } with DAG('my_dag', default_args=default_args, schedule_interval='@daily') as dag: dummy_task = DummyOperator(task_id='dummy_task')
额外小提示
- 调试时可以在回调函数里加日志打印,比如
context['ti'].log.info(f"收件人列表:{email_recipients}"),方便确认变量是否正确获取。 - Airflow中支持Jinja模板的参数,官方文档会标注为"templated",使用前可以先查下参数说明。
内容的提问来源于stack exchange,提问作者Gabe
相关产品推荐
相关产品推荐

