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

MWAA环境下Airflow Redshift重复检测任务实现咨询

MWAA + Redshift 循环检查表重复订单ID并告警的最优实现

针对你的需求,结合MWAA(托管式Airflow)和Redshift环境,推荐用Airflow 2.x的动态任务生成+RedshiftHook+邮件告警方案,解决参数传递和循环遍历问题,具体实现如下:

核心逻辑

  1. 维护需要检查的表列表,避免硬编码
  2. 针对每个表动态生成独立检查任务,支持并行执行
  3. 用RedshiftHook直接连接Redshift执行查询,无需手动管理连接
  4. 根据查询结果判断是否触发邮件告警,可选添加任务汇总

完整DAG代码示例

from airflow import DAG
from airflow.providers.amazon.aws.hooks.redshift import RedshiftHook
from airflow.operators.python import PythonOperator
from airflow.utils.task_group import TaskGroup
from airflow.utils.email import send_email
from datetime import datetime, timedelta

# DAG基础配置
default_args = {
    'owner': 'data-team',
    'depends_on_past': False,
    'email_on_failure': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

# 配置项:按需修改
TABLE_LIST = ['sales.order_table', 'sales.return_order', 'crm.customer_order']  # 要检查的表(含schema)
REDSHIFT_CONN_ID = 'mwaa-redshift-connection'  # MWAA中配置的Redshift连接ID
ALERT_RECEIVERS = ['data-alert@yourcompany.com']  # 告警邮件接收人

def check_duplicate_orders(table_name):
    """检查指定表中的重复订单ID,大于0则发告警邮件"""
    # 初始化Redshift连接钩子
    redshift_hook = RedshiftHook(redshift_conn_id=REDSHIFT_CONN_ID)
    
    # 查询重复订单ID的统计逻辑:统计有多少个重复的order_id组
    check_sql = f"""
        SELECT COUNT(*) AS duplicate_group_count
        FROM (
            SELECT order_id
            FROM {table_name}
            WHERE order_id IS NOT NULL
            GROUP BY order_id
            HAVING COUNT(*) > 1
        ) AS duplicate_groups
    """
    
    # 执行查询并获取结果
    query_result = redshift_hook.get_first(check_sql)
    duplicate_count = query_result[0] if query_result else 0
    
    # 触发告警
    if duplicate_count > 0:
        email_subject = f"【数据质量告警】{table_name} 存在重复订单ID"
        email_body = f"""
        <h3>表 {table_name} 检测到重复订单ID</h3>
        <p>重复的订单ID组数:{duplicate_count}</p>
        <p>查询SQL:</p>
        <pre>{check_sql}</pre>
        """
        send_email(to=ALERT_RECEIVERS, subject=email_subject, html_content=email_body)
    
    return f"表 {table_name} 检查完成,重复组数:{duplicate_count}"

def summarize_check_results(**context):
    """汇总所有表的检查结果,发送汇总邮件"""
    dag_run = context['dag_run']
    summary_content = []
    
    for task_instance in dag_run.get_task_instances():
        if 'check_duplicate_order_ids' in task_instance.task_id:
            table_name = task_instance.task_id.split('__')[1]
            check_result = task_instance.xcom_pull()
            summary_content.append(f"- {table_name}: {check_result}")
    
    summary_email_subject = "重复订单ID检查任务执行汇总"
    summary_email_body = f"""
    <h3>所有表检查结果汇总</h3>
    <ul>
        {''.join([f'<li>{line}</li>' for line in summary_content])}
    </ul>
    """
    send_email(to=ALERT_RECEIVERS, subject=summary_email_subject, html_content=summary_email_body)

with DAG(
    dag_id='redshift_duplicate_order_id_monitor',
    default_args=default_args,
    description='循环检查Redshift表中的重复订单ID并发送告警',
    schedule_interval='0 1 * * *',  # 每天凌晨1点执行
    start_date=datetime(2024, 1, 1),
    catchup=False,
    tags=['redshift', 'mwaa', 'data-quality'],
) as dag:

    # 用TaskGroup管理所有检查任务,保持DAG结构清晰
    with TaskGroup('duplicate_check_tasks') as check_group:
        # 用expand动态生成每个表的检查任务,支持并行执行
        check_task = PythonOperator.partial(
            task_id='check_duplicate_order_ids',
            python_callable=check_duplicate_orders,
        ).expand(op_kwargs=[{'table_name': table} for table in TABLE_LIST])

    # 汇总任务:依赖所有检查任务完成后执行
    summary_task = PythonOperator(
        task_id='summarize_checks',
        python_callable=summarize_check_results,
        provide_context=True,
    )

    check_group >> summary_task

关键配置与注意事项

  • Redshift连接配置:在MWAA的Airflow UI中进入「Admin → Connections」,创建类型为Amazon Redshift的连接,填写主机、端口、数据库名、用户名、密码,确保连接ID与代码中的REDSHIFT_CONN_ID一致。
  • 邮件服务配置:MWAA需要通过环境变量配置SMTP服务,例如设置AIRFLOW__SMTP__SMTP_HOST、AIRFLOW__SMTP__SMTP_PORT、AIRFLOW__SMTP__SMTP_USER等参数,确保send_email能正常发送邮件。
  • SQL安全:因为TABLE_LIST是内部维护的可控列表,所以用format拼接表名是安全的;如果表名来自外部输入,必须添加合法性校验(如正则匹配表名格式),避免SQL注入风险。
  • 动态任务优势:使用expand方法是Airflow 2.x的原生特性,比传统循环创建任务更简洁,且支持并行执行,适配MWAA的分布式运行环境。
  • 错误告警:可以在default_args中添加on_failure_callback,实现任务失败时自动发送告警邮件,示例:
    def task_failure_notify(context):
        fail_subject = f"任务失败:{context['ti'].task_id}"
        fail_body = f"任务 {context['ti'].task_id} 执行失败,错误信息:{str(context['exception'])}"
        send_email(to=ALERT_RECEIVERS, subject=fail_subject, html_content=fail_body)
    
    # 添加到default_args
    default_args['on_failure_callback'] = task_failure_notify
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 16:12:03