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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 06:49:05