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

如何从Airflow变量渲染默认收件人邮件地址列表?

解决方案

方案1:自定义失败回调函数(推荐)

利用Airflow的on_failure_callback参数,在任务失败时主动从Airflow变量中读取收件人列表并发送邮件,绕开email参数不支持模板的限制。

实现步骤:

  • 编写失败回调函数,在函数内读取Airflow变量中的邮件列表,调用Airflow内置邮件工具发送通知
  • 在DAG的default_args中配置on_failure_callback指向该函数

代码示例:

from airflow.models import Variable
from airflow.utils.email import send_email
from datetime import datetime

def task_failure_notify(context):
    # 实时读取Airflow变量中的邮件列表
    email_recipients = Variable.get("my_email_list", deserialize_json=True, default_var=[])
    if not email_recipients:
        return
    
    task_instance = context["task_instance"]
    # 构造邮件内容
    email_subject = f"Airflow任务失败: {task_instance.dag_id}.{task_instance.task_id}"
    email_content = f"""
    <h3>任务失败告警</h3>
    <p>DAG ID: {task_instance.dag_id}</p>
    <p>任务ID: {task_instance.task_id}</p>
    <p>执行时间: {task_instance.execution_date}</p>
    <p>日志链接: <a href="{task_instance.log_url}">点击查看详情</a></p>
    """
    # 发送邮件
    send_email(to=email_recipients, subject=email_subject, html_content=email_content)

# 配置DAG
dag = DAG(
    dag_id="sample_dag",
    default_args={
        "owner": "dev",
        "start_date": datetime(2024, 1, 1),
        "on_failure_callback": task_failure_notify  # 绑定失败回调
    },
    schedule_interval="@daily"
)

方案2:封装通用Operator子类

如果DAG中使用的算子类型相对统一,可以封装带有失败邮件逻辑的Operator子类,所有任务复用该子类,避免重复配置。

代码示例:

from airflow.models import Variable
from airflow.utils.email import send_email
from airflow.operators.python import PythonOperator

class NotifiedPythonOperator(PythonOperator):
    def __init__(self, **kwargs):
        super().__init__(**kwargs)
    
    def on_failure_callback(self, context):
        email_recipients = Variable.get("my_email_list", deserialize_json=True, default_var=[])
        if not email_recipients:
            return
        
        task_instance = context["task_instance"]
        email_subject = f"Airflow任务失败: {task_instance.dag_id}.{task_instance.task_id}"
        email_content = f"""
        <h3>任务失败告警</h3>
        <p>DAG ID: {task_instance.dag_id}</p>
        <p>任务ID: {task_instance.task_id}</p>
        <p>执行时间: {task_instance.execution_date}</p>
        <p>日志链接: <a href="{task_instance.log_url}">点击查看详情</a></p>
        """
        send_email(to=email_recipients, subject=email_subject, html_content=email_content)

# 在DAG中使用自定义算子
with dag:
    sample_task = NotifiedPythonOperator(
        task_id="sample_task",
        python_callable=your_target_function
    )

注意事项:

  • 确保Airflow的SMTP配置已正确完成,否则send_email函数无法正常工作
  • 推荐在回调函数中实时读取Airflow变量,这样邮件列表更新后无需重启DAG即可生效

内容的提问来源于stack exchange,提问作者Gevorg Davoian

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 19:10:32