DAG vững chắc với xử lý lỗi và thử lại

Xây dựng Data Pipeline với Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Vì sao tác vụ thất bại

 

Sự cố kết nối tác vụ

 

  • API timeout và rate limit
  • Mất kết nối ngắn, trễ phân giải tên miền
  • Nhiều tác vụ tranh chấp cùng tài nguyên
Xây dựng Data Pipeline với Airflow

Thử lại (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()

 

  • Thử lại tối đa 3 lần khi lỗi
  • Chờ 2 phút giữa mỗi lần thử
Xây dựng Data Pipeline với Airflow

Exponential backoff

 

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

Backoff lũy thừa

 

  • Độ trễ tăng gấp đôi sau mỗi lần thử lại
  • Cho dịch vụ bên ngoài thời gian phục hồi
Xây dựng Data Pipeline với 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
  • Chuyển tiếp đến Slack, PagerDuty, hoặc công cụ cảnh báo bất kỳ
Xây dựng Data Pipeline với Airflow

Callback ở cấp Dag vs cấp tác vụ

Cấp Dag (tổng quát)

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

Cấp tác vụ (cụ thể)

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

 

  • Cấp tác vụ ghi đè cấp Dag
  • Dùng cả hai để phân tầng cảnh báo
Xây dựng Data Pipeline với Airflow

max_consecutive_failed_dag_runs

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

Các lần chạy Dag thất bại dẫn đến tự tạm dừng

  • Ngăn bội thực cảnh báo do lỗi kéo dài
  • Dag tiếp tục khi được bật lại thủ công
Xây dựng Data Pipeline với Airflow

Kết hợp tất cả

@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 xử lý lỗi tạm thời
  • Callback cảnh báo đúng người
  • Tự tạm dừng giảm tiếng ồn
Xây dựng Data Pipeline với Airflow

Cùng luyện tập nào!

Xây dựng Data Pipeline với Airflow

Preparing Video For Download...