如何使用Apache Airflow执行SQL查询并将结果邮件发送?
简单实现Airflow查询数据表并发送邮件
前置准备
- 在Airflow UI的Admin -> Connections中配置数据库连接,Conn ID设为
my_db_conn,填写对应数据库的主机、端口、账号、密码等信息。 - 修改
airflow.cfg中的SMTP配置,确保Airflow能正常发送邮件(配置smtp_host、smtp_port、smtp_user、smtp_password、smtp_mail_from)。
完整DAG代码
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.hooks.base import BaseHook from airflow.utils.email import send_email from datetime import datetime, timedelta import pandas as pd # DAG默认参数 default_args = { 'owner': 'airflow', 'depends_on_past': False, 'email_on_failure': True, 'email_on_retry': False, 'retries': 1, 'retry_delay': timedelta(minutes=5), } def execute_query_and_send_email(): # 获取预配置的数据库连接 db_conn = BaseHook.get_connection('my_db_conn') # 构造数据库连接字符串(以PostgreSQL为例,其他数据库需调整格式) conn_str = f"postgresql://{db_conn.login}:{db_conn.password}@{db_conn.host}:{db_conn.port}/{db_conn.schema}" # 替换为你的复杂SELECT查询语句 sql_query = """ select * from table where somevar="some_value" """ # 执行查询并转为DataFrame query_result = pd.read_sql(sql_query, conn_str) # 格式化邮件内容 if query_result.empty: email_body = "<p>本次查询无结果返回</p>" else: # 转为HTML表格提升可读性 email_body = query_result.to_html(index=False, border=1) # 发送邮件 send_email( to=['your_email@example.com'], # 替换为你的邮箱 subject='Airflow 数据表查询结果', html_content=f"<h2>查询结果</h2>{email_body}" ) # 定义DAG with DAG( dag_id='db_query_email_dag', default_args=default_args, description='查询指定数据表并发送结果至邮箱', schedule_interval=timedelta(days=1), # 按需调整调度频率 start_date=datetime(2024, 1, 1), catchup=False, ) as dag: # 定义唯一任务 query_email_task = PythonOperator( task_id='run_query_and_send_email', python_callable=execute_query_and_send_email, ) query_email_task
关键注意事项
- 依赖安装:需在Airflow运行环境中安装对应数据库驱动(如PostgreSQL用
pip install psycopg2-binary,MySQL用pip install mysql-connector-python)和pandas。 - 查询适配:根据你的数据库类型调整连接字符串格式,SQL语句直接替换即可,无需额外封装。
- 邮件优化:如果查询结果行数过多,可以添加分页或限制返回行数的逻辑,避免邮件过大。
内容的提问来源于stack exchange,提问作者Probosckie
相关产品推荐
相关产品推荐

