始める無料で始める

コールバック関数の実装

あなたは最近、チームが作成した DAG に失敗時のコールバックを追加する担当になりました。まずは、sales_etl_dag が失敗したときに監査ログへメッセージを書き込む、シンプルな失敗コールバックを追加しましょう。

dagtask オブジェクトはすでにインポートされており、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()
コードを編集して実行