コールバック関数の実装
あなたは最近、チームが作成した DAG に失敗時のコールバックを追加する担当になりました。まずは、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()