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

Cloud Composer配置邮件通知时Broken DAG无法导入email错误如何解决

错误修正方案

核心原因

报错cannot import name 'email'是由Airflow版本迭代导致的导入路径变更引发的,你参考的旧示例代码适配的是低版本Airflow,而当前你使用的Cloud Composer环境搭载的Airflow版本已调整了EmailOperator的归属路径。

具体修改点

  • 替换错误的导入语句
    删除原代码中的from airflow.operators import email,替换为:
from airflow.operators.email import EmailOperator

同时把实例化EmailOperator的代码同步调整,去掉前面的email.前缀:

email_summary = EmailOperator(
    task_id='email_summary',
    to=models.Variable.get('email'),
    subject='testing email service of gcp',
    html_content="""
    Analyzed the Requirements and have created a test POC.Lets see If it works out!
    """
)
  • 优化start_date配置
    你已经导入了days_ago工具,官方不推荐使用动态计算的datetime.now()作为start_date,容易引发调度异常,可将yesterday的赋值修改为:
yesterday = days_ago(1)
  • 补充任务依赖
    当前两个任务没有定义上下游关系,需要在DAG代码末尾添加依赖,实现跑完Dataproc工作流后再发通知邮件:
start_template_job >> email_summary

修正后完整代码

import datetime

from airflow import models
from airflow.contrib.operators import dataproc_operator
from airflow.utils.dates import days_ago
from airflow.operators.email import EmailOperator

project_id = models.Variable.get("project_id")
yesterday = days_ago(1)

default_args = {
    'owner': 'airflow',
    'start_date': yesterday,
    'depends_on_past': False,
    'email': models.Variable.get("email"),
    'email_on_failure': True,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': datetime.timedelta(minutes=5),
}

with models.DAG(
    "dataproc_workflow_dag",
    default_args=default_args,
    schedule_interval=datetime.timedelta(days=1),
) as dag:

    start_template_job = dataproc_operator.DataprocWorkflowTemplateInstantiateOperator(
        task_id="dataproc_workflow_dag",
        template_id="mywf1",
        project_id=project_id,
        region="us-central1",
    )
    email_summary = EmailOperator(
        task_id='email_summary',
        to=models.Variable.get('email'),
        subject='testing email service of gcp',
        html_content="""
        Analyzed the Requirements and have created a test POC.Lets see If it works out!
        """
    )

    # 定义任务上下游
    start_template_job >> email_summary

额外说明

你在default_args中已经配置了email_on_failure=True,只要SendGrid配置正确,流水线任意任务运行失败时,Airflow会自动向配置的邮箱发送失败告警,无需额外编写失败分支逻辑。

内容的提问来源于stack exchange,提问作者Snehil Singh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 21:24:03