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

如何在MWAA工作流中使用Secret Manager存储SMTP信息并发送邮件?

在MWAA中使用Secret Manager存储SMTP信息实现邮件发送

1. 在Secret Manager创建SMTP配置密钥

  • 登录AWS控制台进入Secret Manager服务
  • 创建新密钥,选择「其他类型的密钥」
  • 密钥名称建议用有意义的命名(比如mwaa/smtp-config),密钥内容用JSON格式存储核心SMTP信息:
    {
      "smtp_host": "smtp.example.com",
      "smtp_port": 587,
      "smtp_user": "your-smtp-username@example.com",
      "smtp_password": "your-smtp-password",
      "sender_email": "your-sender-email@example.com"
    }
    
  • 完成创建后,记录密钥的ARN或名称,后续DAG会用到

2. 给MWAA执行角色添加Secret Manager访问权限

  • 进入IAM控制台,找到你的MWAA环境对应的执行角色(通常命名格式为airflow-MwaaEnvironmentName-ExecutionRole-xxxx)
  • 为该角色附加自定义策略,允许读取刚才创建的SMTP密钥:
    {
      "Version": "2012-10-17",
      "Statement": [
          {
              "Effect": "Allow",
              "Action": "secretsmanager:GetSecretValue",
              "Resource": "arn:aws:secretsmanager:your-region:your-account-id:secret:mwaa/smtp-config-xxxx"
          }
      ]
    }
    
    注意替换your-region、your-account-id和密钥的完整ARN(可从Secret Manager控制台复制)

3. DAG中实现邮件发送(两种方式)

方式一:用PythonOperator直接调用smtplib(推荐,更灵活)

不依赖Airflow默认SMTP配置,直接读取密钥后发送邮件:

import boto3
import json
import smtplib
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime

def get_smtp_secret(secret_name):
    """从Secret Manager读取SMTP配置"""
    client = boto3.client('secretsmanager')
    response = client.get_secret_value(SecretId=secret_name)
    return json.loads(response['SecretString'])

def send_email(**context):
    """发送邮件核心逻辑"""
    # 读取SMTP配置
    smtp_config = get_smtp_secret("mwaa/smtp-config")
    
    # 提取配置参数
    smtp_host = smtp_config['smtp_host']
    smtp_port = smtp_config['smtp_port']
    smtp_user = smtp_config['smtp_user']
    smtp_password = smtp_config['smtp_password']
    sender = smtp_config['sender_email']
    
    # 从DAG参数获取收件人、主题和内容
    recipient = context['params']['recipient']
    subject = context['params']['subject']
    body = context['params']['body']
    
    # 构造邮件
    message = f"Subject: {subject}\n\n{body}"
    
    # 发送邮件(如果用465端口,替换为smtplib.SMTP_SSL并移除starttls)
    with smtplib.SMTP(smtp_host, smtp_port) as server:
        server.starttls()
        server.login(smtp_user, smtp_password)
        server.sendmail(sender, recipient, message)

# 定义DAG
default_args = {
    'owner': 'airflow',
    'start_date': datetime(2024, 1, 1),
    'retries': 1
}

with DAG(
    'mwaa_send_email_secret_manager',
    default_args=default_args,
    schedule_interval=None,
    catchup=False
) as dag:
    send_email_task = PythonOperator(
        task_id='send_email_task',
        python_callable=send_email,
        params={
            'recipient': 'target@example.com',
            'subject': 'MWAA测试邮件',
            'body': '这是通过Secret Manager存储SMTP信息发送的测试邮件'
        }
    )

send_email_task

方式二:结合EmailOperator使用

如果偏好Airflow官方的EmailOperator,可通过XCom传递密钥参数:

import boto3
import json
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.email import EmailOperator
from datetime import datetime

def get_smtp_secret(**context):
    client = boto3.client('secretsmanager')
    response = client.get_secret_value(SecretId="mwaa/smtp-config")
    smtp_config = json.loads(response['SecretString'])
    # 将配置推送到XCom供后续任务调用
    context['task_instance'].xcom_push(key='smtp_config', value=smtp_config)

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2024, 1, 1),
    'retries': 1
}

with DAG(
    'mwaa_email_operator_secret',
    default_args=default_args,
    schedule_interval=None,
    catchup=False
) as dag:
    fetch_smtp_config = PythonOperator(
        task_id='fetch_smtp_config',
        python_callable=get_smtp_secret
    )

    send_email = EmailOperator(
        task_id='send_email',
        to='target@example.com',
        subject='EmailOperator测试邮件',
        html_content='<p>这是用EmailOperator结合Secret Manager发送的邮件</p>',
        # 从XCom读取SMTP参数
        smtp_host="{{ task_instance.xcom_pull(task_ids='fetch_smtp_config', key='smtp_config')['smtp_host'] }}",
        smtp_port="{{ task_instance.xcom_pull(task_ids='fetch_smtp_config', key='smtp_config')['smtp_port'] }}",
        smtp_user="{{ task_instance.xcom_pull(task_ids='fetch_smtp_config', key='smtp_config')['smtp_user'] }}",
        smtp_password="{{ task_instance.xcom_pull(task_ids='fetch_smtp_config', key='smtp_config')['smtp_password'] }}",
        from_email="{{ task_instance.xcom_pull(task_ids='fetch_smtp_config', key='smtp_config')['sender_email'] }}",
        smtp_starttls=True
    )

fetch_smtp_config >> send_email

注意事项

  • MWAA环境默认预装boto3,无需额外安装依赖
  • 若SMTP服务器使用465端口,需将smtplib.SMTP替换为smtplib.SMTP_SSL,并移除server.starttls()调用
  • 遵循最小权限原则,仅给MWAA执行角色授予必要的Secret Manager访问权限
  • 测试时可通过CloudWatch日志排查邮件发送失败问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 16:30:30