實作回呼函式
你最近被指派要替團隊建立的 Dags 加上失敗回呼。首先,你想在 sales_etl_dag 失敗時,寫一則訊息到稽核紀錄中,作為一個簡單的失敗回呼。
dag 與 task 物件已經匯入,get_sales_data 與 process_sales_data 任務也都已建立。
本練習屬於課程
Python 中的 Apache Airflow 入門
練習說明
- 建立一個名為
alert_on_failure的回呼函式。 - 定義該函式可接受 Airflow 傳入的任何物件。
- 使用
alert_on_failure函式來指定失敗回呼。
動手互動練習
試著完成這個範例程式碼,體驗一下這個練習。
# Create the callback function
def ____(____):
dag_id = context["dag"].dag_id
task_id = context["task_instance"].task_id
print(f"Task {task_id} in Dag {dag_id} has failed.")
# Specify the Dag with a failure callback
@dag(dag_id='sales_etl_dag',
____=alert_on_failure
)
def sales_etl_dag():
get_sales_data() >> process_sales_data()
sales_etl_dag()