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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 09:15:37