如何在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
相关产品推荐
相关产品推荐

