开始使用免费开始使用

实现回调函数

您最近被指派为团队创建的多个 Dag 添加失败回调。作为开始,您希望在 sales_etl_dag 失败时,添加一个简单的失败回调,将消息写入审计日志。

dagtask 对象已导入,get_sales_dataprocess_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()
编辑并运行代码