实现 Trigger Rule
创建工作流后,您意识到如果至少有一个任务失败,Dag 能提供一些状态更新会更好。您决定实现一个任务,在 Dag 上执行 one failed 检查,以便当 Dag 中任一任务失败时提醒您。
其他所有任务都已定义,且 task 与 dag 对象已为您导入。
本练习是课程的一部分
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()