การสร้าง Callback Function
คุณได้รับมอบหมายให้เพิ่ม failure callback ให้กับ DAG ที่ทีมสร้างขึ้น โดยเริ่มต้นจากการเพิ่ม failure callback อย่างง่าย ซึ่งจะบันทึกข้อความลงใน audit log เมื่อ sales_etl_dag เกิดข้อผิดพลาด
ออบเจกต์ dag และ task ถูก import ไว้แล้ว และ task get_sales_data กับ process_sales_data ถูกสร้างขึ้นแล้วเช่นกัน
แบบฝึกหัดนี้เป็นส่วนหนึ่งของหลักสูตร
Apache Airflow เบื้องต้นด้วย Python
คำแนะนำการฝึกหัด
- สร้าง callback function ชื่อ
alert_on_failure - กำหนดให้ฟังก์ชันรับออบเจกต์ที่ Airflow ส่งมาได้ทุกรูปแบบ
- ระบุ failure callback โดยใช้ฟังก์ชัน
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()