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

如何在DAG-B中获取DAG-A的最新运行状态并传入Slack消息

跨DAG获取最新运行状态并集成到Slack通知

核心思路

直接通过Airflow的元数据模型查询目标DAG(DAG-A)的运行记录,获取最新状态,全程无需修改DAG-A的代码,所有逻辑在DAG-B中实现。

实现步骤

1. 编写获取DAG-A最新状态的函数

利用Airflow内置的DagRun模型查询指定DAG的最新运行记录,提取状态并格式化:

from airflow.models import DagRun
from airflow.utils.state import State

def _get_message() -> str:
    # 替换为你的DAG-A的ID
    target_dag_id = "data_pipeline"
    # 按执行时间倒序取第一条运行记录
    latest_dag_run = DagRun.find(dag_id=target_dag_id, order_by=["-execution_date"]).first()
    
    if latest_dag_run:
        status = latest_dag_run.state
        # 给不同状态添加可视化标识
        if status == State.SUCCESS:
            status_text = f"*{status}* ✅"
        elif status == State.FAILED:
            status_text = f"*{status}* ❌"
        else:
            status_text = f"*{status}* ⚠️"
        return f"Status of data_pipeline: {status_text}"
    else:
        return "Status of data_pipeline: 未找到运行记录"

2. 整合到DAG-B的Slack通知任务

将上述函数替换原有的_get_message,直接对接Slack通知任务:

from datetime import datetime
from airflow import DAG
from airflow.providers.slack.operators.slack_webhook import SlackWebhookOperator

default_args = {
    'owner': 'airflow',
    'retries': 1,
}

def _get_message() -> str:
    from airflow.models import DagRun
    from airflow.utils.state import State
    target_dag_id = "data_pipeline"
    latest_dag_run = DagRun.find(dag_id=target_dag_id, order_by=["-execution_date"]).first()
    
    if latest_dag_run:
        status = latest_dag_run.state
        if status == State.SUCCESS:
            status_text = f"*{status}* ✅"
        elif status == State.FAILED:
            status_text = f"*{status}* ❌"
        else:
            status_text = f"*{status}* ⚠️"
        return f"Status of data_pipeline: {status_text}"
    else:
        return "Status of data_pipeline: 未找到运行记录"

with DAG("slack_dag", 
         start_date=datetime(2021, 1, 1), 
         schedule_interval="@daily", 
         default_args=default_args, 
         catchup=False
        ) as dag:

    send_slack_notification = SlackWebhookOperator(
        task_id="send_slack_notification",
        http_conn_id="slack_conn",
        message=_get_message(),
        channel="#test-public"
    )

关键说明

  • 权限要求:运行DAG-B的Airflow Worker默认已具备访问元数据库的权限,无需额外配置。
  • 状态准确性:使用Airflow内置的State枚举类,避免手动拼写状态值导致的错误,涵盖SUCCESS、FAILED、RUNNING等所有官方状态。
  • 性能保障:DagRun.find直接查询元数据,对于每日多次运行的DAG,查询效率足够,无需额外缓存优化。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 15:52:37