实现回调函数
您最近被指派为团队创建的多个 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()