Triển khai Trigger Rule
Sau khi tạo một workflow, bạn nhận ra Dag sẽ hữu ích hơn nếu cung cấp cập nhật khi có ít nhất một task thất bại. Bạn quyết định thêm một task thực hiện kiểm tra one failed trên Dag để cảnh báo bạn nếu bất kỳ task nào trong Dag bị lỗi.
Tất cả các task khác đã được định nghĩa và các đối tượng task & dag đã được nhập sẵn cho bạn.
Bài tập này là một phần của khóa học
Giới thiệu về Apache Airflow bằng Python
Hướng dẫn bài tập
- Import thư viện phù hợp để dùng trigger rules.
- Thêm thuộc tính trigger rule thích hợp cho task
notify_on_failure. - Thiết lập thuộc tính để task được kích hoạt khi một hoặc nhiều task upstream thất bại.
- Đặt
notify_on_failurelà phụ thuộc downstream của hai task transform.
Bài tập tương tác thực hành trực tiếp
Hãy thử làm bài tập này bằng cách hoàn thành đoạn mã mẫu này.
# 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()