เริ่มต้นใช้งานเริ่มต้นใช้งานได้ฟรี

การสร้าง 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()
แก้ไขและรันโค้ด