on_failure_callback를 통한 Dag 수준 재시도

on_failure_callback를 통한 Dag 수준 재시도 (Dag-level Retry via on_failure_callback)

태스크가 재시도를 모두 소진한 뒤 실패했을 때, Dag 실행 전체를 다시 시도하는 패턴을 설명하는 문서예요. on_failure_callback과 공개 REST API의 Clear 엔드포인트를 조합해 실패한 Dag run을 지우고 스케줄러가 다시 실행하게 하는 레시피예요. 재시도 횟수 제한, 한계(caveats), 자동 복구 가드까지 함께 살펴볼게요.

출처: 문서

본문

Airflow의 내장 재시도 메커니즘은 task 수준에서 동작해요: 각 태스크는 자신만의 retriesretry_delay를 갖고, 실패 시 그 단일 태스크만 다시 실행돼요.

태스크는 이상적으로 task 수준에서 멱등(idempotent) 하게 설계되어 내장 태스크별 재시도 메커니즘으로 충분해야 해요. 이것이 권장되는 출발점이고, 대부분의 워크플로우는 그것 이상이 필요하지 않아요.

이 how-to는 태스크 수준 멱등성으로도 더 큰 단위의 재시도 — 사실상 Dag 수준 재시도 — 를 원하는 경우를 위한 것이에요. 예를 들어 중간 상태를 부분적으로 정리하기 어려운 다단계 파이프라인은, 실패한 단일 태스크를 다시 실행하는 것보다 처음부터 다시 실행하는 게 더 안전할 수 있어요.

이 문서는 on_failure_callback과 공개 Airflow REST API를 결합해 Dag 수준 재시도에 근사하는 패턴을 보여줘요: 태스크가 실패하면 콜백이 Clear 엔드포인트를 호출해 실패한 Dag run을 지우고, 스케줄러가 다시 실행하게 해요.

이것은 기존 프리미티브 위에 구축된 레시피이지 새 기능이 아니에요. 의견이 반영되어 있고 트레이드오프가 있어요. 사용하기 전에 아래 '한계(Caveats)'를 주의 깊게 읽어 보세요.

이에 대한 더 일반적인 구조 — 때때로 transactional task group 이라고도 불려요 — 가 Airflow dev list에서 추가 기능으로 논의되고 있어요. 유스케이스에 도움이 된다면, 이 레시피를 장기적으로 의존하기보다 그 논의에 참여하세요.

이 패턴이 언제 적합할까요?

이런 경우에 사용하세요:

  • 작업의 단위가 자연스럽게 단일 태스크가 아니라 Dag run일 때.
  • 태스크 재실행이 멱등할 때 — 다시 실행해도 중복된 부작용, 이중 청구, 불일치한 외부 상태가 생기지 않을 때.
  • 시간과 비용 측면에서 제한된 횟수의 전체 재실행이 허용 가능할 때.

이런 경우엔 피하세요:

  • 태스크에 비멱등 부작용(이메일 보내기, 결제 청구, 중복 제거되지 않는 API에 게시)이 있고 멱등하게 만들지 않았을 때.
  • 재시도가 다운스트림 Dag나 asset에 투명해야 할 때 — Dag run을 지우면 원래 실행의 재시도가 아니라 새로운 시도가 생겨요.
  • 실제로 보이는 실패 모드를 단순한 태스크별 retries 설정이 이미 커버할 때.

동작 방식 (How it works)

  1. 태스크가 태스크 수준 재시도를 모두 소진한 뒤 실패해요.
  2. 태스크의 on_failure_callback이 dag processor에서 실행돼요.
  3. 콜백이 (자신이 유지하는 카운터를 사용해 — 아래 '재시도 횟수 제한' 참고) 또 다른 Dag 수준 시도가 허용되는지 판단해요.
  4. 허용된다면 콜백이 Airflow REST API를 호출해 실패한 Dag run을 지워요.
  5. 스케줄러가 지워진 task instance들을 집어 다시 실행해요.

사용되는 REST 엔드포인트는 POST /api/v2/dags/{dag_id}/dagRuns/{dag_run_id}/clear예요.

전제 조건 (Prerequisites)

  • dag processor에서 Airflow REST API로의 인증된 접근, 대상 Dag의 Dag run을 지울 권한.
  • dag processor에서 API 서버로의 네트워크 도달성.

참고 (Note)

on_failure_callback은 정상적인 태스크 실행을 통해 실패가 발생할 때만 실행돼요. 수동으로 수행한 상태 변경(UI, CLI)은 이 콜백을 호출하지 않아요.

기본 예시: 실패한 Dag run 지우기 (Basic example: clear the failed Dag run)

다음 예시는 태스크 수준 재시도를 소진한 뒤 Dag run 내의 어떤 태스크든 실패하면 전체 Dag run을 지워요. 콜백이 첫 실패에 실행되고 전체 Dag run이 재시도 책임을 넘겨받도록 태스크에 retries=0을 설정했어요.

이 최소 버전은 실패 시 항상 지워요. 프로덕션에서 그대로 실행하지 마세요 — 태스크가 계속 실패하면 무한 루프에 빠질 수 있어요. 배포 전에 시도 횟수 제한을 추가하세요. 아래 '재시도 횟수 제한'을 참고하세요.

from urllib.parse import quote

import requests

from airflow.sdk import DAG, task

AIRFLOW_API_BASE = "https://airflow.example.com/api/v2"  # your deployment's API base

def clear_dag_run_on_failure(context):
    dag_run = context["dag_run"]
    dag_id_path = quote(dag_run.dag_id, safe="")
    run_id_path = quote(dag_run.run_id, safe="")
    response = requests.post(
        f"{AIRFLOW_API_BASE}/dags/{dag_id_path}/dagRuns/{run_id_path}/clear",
        headers={"Authorization": "Bearer <token>"},
        json={"dry_run": False},
        timeout=30,
    )
    response.raise_for_status()

with DAG(
    dag_id="example_dag_level_retry",
    default_args={
        "retries": 0,  # let the Dag-level retry take over
        "on_failure_callback": clear_dag_run_on_failure,
    },
):

    @task
    def step_one(): ...

    @task
    def step_two(value): ...

    step_two(step_one())

재시도 횟수 제한하기 (Limiting the number of retries)

항상 run을 지우는 순진한 콜백은 태스크가 계속 실패하면 무한 루프를 돌 거예요. dag_run.conf는 run이 생성될 때 설정되고 clear에 의해 변경되지 않으므로 재시도 카운터로 사용할 수 없어요 — 시도 횟수를 직접 추적해야 해요. 간단한 방법은 run에 스코프된 Variable을 사용하는 거예요:

from airflow.sdk import Variable

def _attempts_key(dag_id: str, run_id: str) -> str:
    return f"dag_run_attempts::{dag_id}::{run_id}"

def _read_attempts(dag_id: str, run_id: str) -> int:
    return int(Variable.get(_attempts_key(dag_id, run_id), default=0))

def _bump_attempts(dag_id: str, run_id: str) -> int:
    new_value = _read_attempts(dag_id, run_id) + 1
    Variable.set(_attempts_key(dag_id, run_id), str(new_value))
    return new_value

지우기로 결정하기 전에 _read_attempts를 사용하고, API 호출 직전에 _bump_attempts를 호출해요. Dag run이 끝날 때(성공 시와 포기 시 모두) 카운터를 리셋해서 Variables가 쌓이지 않게 해요:

def cleanup_attempts(context):
    dag_run = context["dag_run"]
    Variable.delete(_attempts_key(dag_run.dag_id, dag_run.run_id))

cleanup_attempts를 Dag 수준의 on_success_callback에 연결하고, on_failure_callback의 포기(give-up) 분기에서 반환 전에 호출해요.

참고 (Note)

이 Variable 기반 카운터는 동시 콜백 간에 원자적(atomic)이지 않아요. Dag에 동시에 실패할 수 있는 병렬 분기가 있다면, 두 콜백이 같은 값을 읽고 둘 다 지우기로 결정할 수 있어요. 실제로는 두 번째 clear가 no-op(first가 이미 run을 리셋)이지만 중복 API 호출이 보일 수 있어요. 더 엄격한 보증이 필요하다면 compare-and-swap 시맨틱이 있는 백엔드를 사용하세요.

자동 복구 Dags 가드하기 (Guarding automated remediation Dags)

Dag 수준 재시도 패턴은 실패한 작업을 지우거나, 메시지를 재생하거나, 외부 job을 재시작하는 등 외부 시스템에 영향을 주는 작업을 실행하는 자동 복구 Dags에서도 사용할 수 있어요. 이러한 작업은 지속적인 실패가 Dag 실행을 가로질러 같은 복구 작업을 반복 실행하게 하지 않도록 가드레일(guardrail)이 있어야 해요.

Pool은 동시에 실행되는 복구 태스크 수를 제한할 수 있지만, 같은 복구 작업이 시도되는 빈도는 제한하지 않아요. 자동 복구 워크플로우에서는 각 대상에 작은 cooldown 마커와 수동 검토 전 최대 자동 작업 횟수를 추가하는 것을 고려해 보세요.

cooldown 타임스탬프나 시도 카운터 같은 작은 가드 값에는 Airflow Variable로 충분할 수 있어요. Variables를 대용량 상태 저장소나 강한 일관성을 가진 잠금 메커니즘으로 사용하지 마세요. 워크플로우가 원자적 업데이트, 엄격한 동시성 제어, 또는 대상별 상태 레코드가 많아야 한다면 그 목적에 맞게 설계된 외부 저장소를 사용하세요.

한계 (Caveats)

  • **Dag 수준 재시도 사이의 UP_FOR_RETRY 대기. ** 실패한 태스크는 다음 clear를 트리거하기 전에 UP_FOR_RETRY에서 retry_delay(기본 5분) 동안 대기한 뒤 FAILED로 전환되므로, 각 Dag 수준 재시도는 최소 그 간격만큼 지연돼요.
  • **멱등성은 여러분의 책임이에요. ** 지워진 모든 태스크가 다시 실행돼요. 실패한 시도에서 외부 부작용을 만들었다면 — 행(row) 쓰기, 메시지 전송, 파일 업로드 — 태스크가 멱등하지 않다면(예: upsert, 중복 제거 키, 조건부 쓰기 사용) 다시 생성돼요.
  • **루프 위험. ** 카운터나 실패하는 태스크의 버그가 무한 clear-and-retry 루프를 일으킬 수 있어요. 항상 시도 횟수에 상한을 두고 상한에 도달하면 알림을 받으세요.
  • **비용. ** Dag 수준 재시도는 지워진 모든 태스크를 처음부터 다시 실행해요. 시간과 리소스 비용이 허용 가능한지 확인하고, 시도 상한을 낮게 유지하세요.
  • **실행 중인 태스크. ** Dag에 병렬 분기가 있다면, 한 태스크가 실패하는 동안 같은 Dag run의 다른 태스크가 아직 실행 중일 수 있어요. 다른 작업이 진행 중일 때 run을 지우는 것이 안전하도록 Dag와 그 부작용을 설계하세요.
  • **토큰 수명. ** 대부분의 auth manager가 발급한 토큰은 만료돼요(예: simple auth manager의 JWT는 기본 24시간). Variable에 저장된 정적 토큰은 만료 후 조용히 401을 반환해요. 루프 가드가 잘못된 동작을 막지만 Dag 수준 재시도는 동작을 멈춰요. 장기 실행 배포에서는 자격 증명을 주기적으로 갱신하거나 더 긴 수명의 서비스 자격 증명을 사용하세요.

더 알아보기 (Learn more)