如何配置Cloud Composer以实现邮件通知功能?
Hey there! Let's break down exactly how to configure your Cloud Composer environment to send emails and get notifications—since Cloud Composer is built on Apache Airflow, we'll leverage Airflow's native email tools to make this work smoothly.
Cloud Composer manages Airflow's configuration via environment variables, so we need to set up the SMTP credentials your environment will use to send emails. Here are the key variables you'll need to add:
AIRFLOW__SMTP__SMTP_HOST: Your SMTP server (e.g.,smtp.gmail.comfor Gmail,smtp.office365.comfor Outlook)AIRFLOW__SMTP__SMTP_PORT: Typically587for TLS connections (most common) or465for SSLAIRFLOW__SMTP__SMTP_USER: The email address you want to send from (e.g.,your-email@gmail.com)AIRFLOW__SMTP__SMTP_PASSWORD: For Gmail/Outlook with 2FA enabled, use an app-specific password (not your regular account password). For other providers, use your account password.AIRFLOW__SMTP__SMTP_MAIL_FROM: Same asSMTP_USER(the "from" address in outgoing emails)AIRFLOW__SMTP__SMTP_STARTTLS: Set toTrueif using port 587AIRFLOW__SMTP__SMTP_SSL: Set toFalseif using STARTTLS,Trueif using port 465
To add these variables:
- Go to your Cloud Composer environment in the Google Cloud Console
- Navigate to the Environment variables tab
- Click Add variable and input each key-value pair one by one
Once SMTP is set up, you can build DAGs that send emails on demand. Here are two common approaches:
Using Airflow's EmailOperator (simple, no custom code)
This is great for basic, pre-defined emails:
from airflow import DAG from airflow.operators.email import EmailOperator from datetime import datetime default_args = { 'owner': 'your-team', 'start_date': datetime(2024, 1, 1), } with DAG( 'daily_status_email', default_args=default_args, schedule_interval='@daily', catchup=False ) as dag: send_status_email = EmailOperator( task_id='send_daily_status', to=['team-member1@example.com', 'team-member2@example.com'], subject='Daily ETL Pipeline Status', html_content='<p>Your daily ETL pipeline ran successfully! Check the Cloud Composer UI for details.</p>' ) send_status_email
Using a PythonOperator (custom, dynamic content)
Use this if you need to generate email content dynamically (e.g., include task metrics or error details):
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.utils.email import send_email from datetime import datetime def send_custom_failure_alert(context): # Pull task details from the execution context task_instance = context['task_instance'] dag_id = task_instance.dag_id task_id = task_instance.task_id error_msg = str(context.get('exception', 'Unknown error')) subject = f"ALERT: Task {task_id} in DAG {dag_id} Failed" html_content = f""" <h3>Task Failure Alert</h3> <p>DAG: {dag_id}</p> <p>Task: {task_id}</p> <p>Error Message:</p> <pre>{error_msg}</pre> """ send_email(to=['on-call-team@example.com'], subject=subject, html_content=html_content) default_args = { 'owner': 'your-team', 'start_date': datetime(2024, 1, 1), } with DAG( 'etl_pipeline_with_alerts', default_args=default_args, schedule_interval='@hourly', catchup=False ) as dag: # Example task that might fail run_etl = PythonOperator( task_id='run_etl_process', python_callable=your_etl_function, on_failure_callback=send_custom_failure_alert # Trigger email on failure ) run_etl
If you want emails sent automatically when tasks succeed or fail, you can configure this directly in your DAG's default_args (applies to all tasks) or per-task:
Global DAG-level notifications
default_args = { 'owner': 'your-team', 'start_date': datetime(2024, 1, 1), 'email': ['team-alerts@example.com'], 'email_on_failure': True, # Send email if any task fails 'email_on_retry': True, # Send email when a task retries 'email_on_success': True # Send email when tasks complete successfully } # Define your DAG and tasks as usual
Per-task notifications (override global settings)
critical_task = PythonOperator( task_id='critical_data_load', python_callable=load_critical_data, # Override global settings for this specific task email_on_failure=True, email_on_success=False, email=['senior-engineer@example.com'] )
- Gmail authentication errors: Make sure you've enabled 2-Step Verification on your Google account and created an App Password (regular passwords won't work for SMTP).
- SMTP connection timeouts: Check that your Cloud Composer environment has outbound network access to your SMTP server (e.g., allow traffic to
smtp.gmail.com:587in your VPC firewall rules). - Emails not showing up: Check the Airflow scheduler logs in Cloud Composer's Logs tab—you'll see error messages if there's an issue with SMTP credentials or server settings.
内容的提问来源于stack exchange,提问作者James

