失敗処理とリトライによる堅牢な DAG

Airflow によるデータパイプラインの構築

Volker Janz

Senior Developer Advocate at Astronomer

タスクが失敗する原因

 

タスクの接続障害

 

  • API のタイムアウトとレート制限
  • 一時的なネットワーク障害やDNS解決の遅延
  • 同じリソースへの複数タスクの競合
Airflow によるデータパイプラインの構築

リトライ

@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()

 

  • 失敗時に最大 3 回 リトライします
  • 各試行の間に 2 分 待機します
Airflow によるデータパイプラインの構築

指数バックオフ

 

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

指数バックオフ

 

  • リトライごとに待機時間が 2 倍 になります
  • 外部サービスが復旧する時間を確保します
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():
    ...

 

  • context ディクショナリ:dagtiexception、ログ URL
  • Slack、PagerDuty などのアラートツールに通知できます
Airflow によるデータパイプラインの構築

DAG レベルとタスクレベルのコールバック

DAG レベル(全タスク対象)

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

タスクレベル(個別指定)

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

 

  • タスクレベルは DAG レベルを上書きします
  • 両方を組み合わせて多段階のアラートを設定できます
Airflow によるデータパイプラインの構築

max_consecutive_failed_dag_runs

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

連続失敗後の自動一時停止

  • 継続的な失敗によるアラート過多を防ぎます
  • 手動で一時停止を解除すると DAG が再開します
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():
        ...
  • リトライ+バックオフで一時的な障害に対処します
  • コールバックで適切な担当者に通知します
  • 自動一時停止で不要なアラートを抑制します
Airflow によるデータパイプラインの構築

練習しましょう!

Airflow によるデータパイプラインの構築

Preparing Video For Download...