DAG ที่แข็งแกร่งด้วยการจัดการความล้มเหลวและการลองใหม่

การสร้าง Data Pipeline ด้วย Airflow

Volker Janz

Senior Developer Advocate at Astronomer

สาเหตุที่ Task ล้มเหลว

 

ปัญหาการเชื่อมต่อของ Task

 

  • API หมดเวลาและถึงขีดจำกัดอัตราการใช้งาน
  • การขัดข้องของเครือข่ายชั่วคราวและความล่าช้าในการแปลงชื่อโดเมน
  • หลาย Task แย่งใช้ทรัพยากรเดียวกัน
การสร้าง Data Pipeline ด้วย Airflow

การลองใหม่ (Retries)

@task(
    retries=3,
    retry_delay=timedelta(minutes=2),
)
def fetch_weather():
    response = requests.get("https://api.weather.com/forecast")
    response.raise_for_status()
    return response.json()

 

  • Task ลองใหม่สูงสุด 3 ครั้ง เมื่อล้มเหลว
  • รอ 2 นาที ระหว่างแต่ละครั้ง
การสร้าง Data Pipeline ด้วย Airflow

Exponential Backoff

 

@task(
    retries=3,
    retry_delay=timedelta(minutes=2),
    retry_exponential_backoff=2.0,
)
def fetch_weather():
    ...

Exponential backoff

 

  • หน่วงเวลาเพิ่มเป็นสองเท่าหลังแต่ละครั้งที่ลองใหม่
  • ให้เวลาบริการภายนอกฟื้นตัว
การสร้าง Data Pipeline ด้วย Airflow

on_failure_callback

def alert_on_failure(context):
    dag_id = context["dag"].dag_id
    task_id = context["ti"].task_id
    print(f"ALERT: {task_id} in {dag_id} failed!")

@task(on_failure_callback=alert_on_failure)
def fetch_weather():
    ...

 

  • dict context: dag, ti, exception, log URL
  • ส่งการแจ้งเตือนไปยัง Slack, PagerDuty หรือเครื่องมืออื่น
การสร้าง Data Pipeline ด้วย Airflow

Callback ระดับ DAG กับระดับ Task

ระดับ DAG (ครอบคลุมทุก Task)

@dag(
    on_failure_callback=alert_team,
)
def my_pipeline():
    ...

ระดับ Task (เฉพาะเจาะจง)

@task(
    on_failure_callback=page_oncall,
)
def critical_step():
    ...

 

  • ระดับ Task แทนที่ระดับ DAG
  • ใช้ทั้งสองระดับเพื่อแจ้งเตือนแบบเป็นชั้น
การสร้าง Data Pipeline ด้วย Airflow

max_consecutive_failed_dag_runs

@dag(max_consecutive_failed_dag_runs=3)
def monitoring_pipeline():
    ...

การรัน DAG ที่ล้มเหลวตามด้วยการหยุดอัตโนมัติ

  • ลดการแจ้งเตือนซ้ำจากความล้มเหลวต่อเนื่อง
  • DAG กลับมาทำงานเมื่อเปิดใช้งานด้วยตนเอง
การสร้าง Data Pipeline ด้วย Airflow

สรุปภาพรวม

@dag(
    max_consecutive_failed_dag_runs=3,
    on_failure_callback=alert_team,
)
def weather_pipeline():

    @task(
        retries=3,
        retry_delay=timedelta(minutes=2),
        retry_exponential_backoff=True,
        on_failure_callback=page_oncall,
    )
    def fetch_weather():
        ...
  • Retries + backoff รับมือกับความล้มเหลวชั่วคราว
  • Callbacks แจ้งเตือนคนที่เกี่ยวข้อง
  • Auto-pause ลดการแจ้งเตือนที่ไม่จำเป็น
การสร้าง Data Pipeline ด้วย Airflow

มาฝึกกันเถอะ!

การสร้าง Data Pipeline ด้วย Airflow

Preparing Video For Download...