Apache Airflow:含Skipped Task时DAG显示Success的告警处理需求
解决Airflow DAG含Skipped Task时仍显示Success并触发Slack告警的方案
Airflow默认逻辑是:只要DAG中没有Failed或Upstream Failed状态的Task,整个DAG就会标记为Success,哪怕存在Skipped的Task。要实现你要的效果,核心是新增一个收尾检查Task,在所有业务Task执行完成后,检查是否存在Skipped任务,若存在则主动触发失败并发送Slack告警。
实现步骤
1. 编写检查Skipped任务的函数
用PythonOperator执行检查逻辑,获取当前DAG运行的所有Task实例,筛选出Skipped状态的任务,存在则触发Slack告警并抛出异常(让当前Task变为Failed,进而使整个DAG状态变为Failed):
from airflow.models import TaskInstance from airflow.utils.state import State from airflow.operators.python import PythonOperator from airflow.providers.slack.operators.slack_webhook import SlackWebhookOperator def check_skipped_tasks(**context): dag_run = context["dag_run"] # 获取当前DAG运行周期内的所有Task实例 task_instances = TaskInstance.find(dag_id=dag_run.dag_id, execution_date=dag_run.execution_date) # 筛选出状态为Skipped的Task ID skipped_task_ids = [ti.task_id for ti in task_instances if ti.state == State.SKIPPED] if skipped_task_ids: skipped_task_str = ", ".join(skipped_task_ids) alert_msg = f"DAG [{dag_run.dag_id}] 执行存在Skipped任务: {skipped_task_str}" # 发送Slack告警 SlackWebhookOperator( task_id="slack_skipped_alert", slack_webhook_conn_id="your_slack_webhook_conn", # 替换为你的Slack连接ID message=alert_msg, username="Airflow告警机器人", icon_emoji=":red_circle:" ).execute(context={}) # 抛出异常,让当前Task变为Failed,触发DAG状态变为Failed raise Exception(alert_msg)
2. 在DAG中集成检查Task
将检查Task添加到DAG末尾,设置trigger_rule='all_done'(确保不管前面的业务Task是Success还是Skipped,检查Task都会执行),并让它依赖所有业务Task:
from airflow import DAG from airflow.providers.databricks.operators.databricks import DatabricksRunNowOperator from datetime import datetime default_args = { "owner": "airflow", "start_date": datetime(2024, 1, 1), } with DAG( dag_id="your_target_dag_id", default_args=default_args, schedule_interval="@daily", catchup=False ) as dag: # 你的Databricks业务Task示例 databricks_job_1 = DatabricksRunNowOperator( task_id="databricks_job_1", databricks_conn_id="your_databricks_conn", # 替换为你的Databricks连接ID job_id=123, # 替换为你的Databricks Job ID ) databricks_job_2 = DatabricksRunNowOperator( task_id="databricks_job_2", databricks_conn_id="your_databricks_conn", job_id=456, ) # 添加检查Skipped任务的Task check_skipped_task = PythonOperator( task_id="check_skipped_tasks", python_callable=check_skipped_tasks, provide_context=True, trigger_rule="all_done", # 关键配置:确保所有业务Task完成后必执行 ) # 设置依赖关系:所有业务Task执行完后执行检查 [databricks_job_1, databricks_job_2] >> check_skipped_task
关键说明
trigger_rule='all_done':必须设置这个规则,否则如果前面的业务Task被Skipped,检查Task会默认继承Skipped状态,无法执行检查逻辑。- DAG状态变更逻辑:只要DAG中有一个Task处于
Failed状态,整个DAG就会被标记为Failed(红色圆圈),符合你的需求。 - Slack告警集成:如果已经有现成的Slack告警函数,可以直接替换掉示例中的
SlackWebhookOperator调用部分。
内容的提问来源于stack exchange,提问作者jonhatan_schilino
相关产品推荐
相关产品推荐

