开始使用免费开始使用

实现 Trigger Rule

创建工作流后,您意识到如果至少有一个任务失败,Dag 能提供一些状态更新会更好。您决定实现一个任务,在 Dag 上执行 one failed 检查,以便当 Dag 中任一任务失败时提醒您。

其他所有任务都已定义,且 taskdag 对象已为您导入。

本练习是课程的一部分

Python 中的 Apache Airflow 入门

查看课程

练习说明

  • 导入用于触发规则的合适库。
  • notify_on_failure 任务添加相应的触发规则属性。
  • 设置该属性,使该任务在一个或多个上游任务失败时触发。
  • notify_on_failure 设为两个 transform 任务的下游依赖。

交互式实操练习

通过完成这段示例代码来试试这个练习。

# Import TriggerRule
from airflow.utils.____ import ____

@dag(schedule="@daily", start_date=datetime(2026, 5, 1))
def etl_pipeline():

    # Trigger notify_on_failure when any upstream task fails
    @task(____=TriggerRule.____)
    def notify_on_failure(**context) -> None:
        dag_id = context["dag"].dag_id
        run_id = context["run_id"]
        print(f"ALERT: A task failed in DAG '{dag_id}', run '{run_id}'. Sending notification...")

    # Set notify_on_failure downstream of both transform tasks
    [transform_users(), transform_orders()] ____ notify_on_failure()
    
etl_pipeline()
编辑并运行代码