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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 16:40:32